"""Redis 客户端 + Redis Key 命名空间。 提供进程级单例的 async Redis 连接池,以及统一的 key 前缀生成。 """ from __future__ import annotations import logging from typing import Optional from redis.asyncio import Redis, from_url from app.config import settings logger = logging.getLogger(__name__) _redis: Optional[Redis] = None def _k(suffix: str) -> str: """构造带前缀的 Redis key。""" return f"{settings.redis_key_prefix}:{suffix}" # 统一管理的 Redis Key KEY_CACHE_DATA = _k("cache:data") # JSON 缓存数据 KEY_LOCK_LEADER = _k("lock:leader") # leader 选举锁(值=持有者标识) KEY_LOCK_REFRESH = _k("lock:refresh") # 防并发刷新锁 async def get_redis() -> Redis: """获取进程级 Redis 单例(懒初始化)。""" global _redis if _redis is None: _redis = from_url( settings.redis_url, username=settings.redis_username, password=settings.redis_password, encoding="utf-8", decode_responses=True, ) logger.info("Redis 连接已建立: %s", settings.redis_url) return _redis async def close_redis() -> None: """关闭 Redis 连接(应用关闭时调用)。""" global _redis if _redis is not None: try: await _redis.aclose() except Exception as exc: # noqa: BLE001 logger.warning("关闭 Redis 连接异常: %s", exc) _redis = None async def try_acquire_lock( key: str, value: str, ttl_seconds: int ) -> bool: """尝试获取分布式锁(SET NX EX)。成功返回 True。""" r = await get_redis() ok = await r.set(key, value, nx=True, ex=ttl_seconds) return bool(ok) async def renew_lock(key: str, value: str, ttl_seconds: int) -> bool: """续期锁(仅当 value 与持有者一致时才续期,避免误续他人锁)。 使用 Lua 脚本保证原子性。 """ r = await get_redis() script = """ if redis.call('get', KEYS[1]) == ARGV[1] then return redis.call('expire', KEYS[1], ARGV[2]) else return 0 end """ result = await r.eval(script, 1, key, value, ttl_seconds) return bool(result) async def release_lock(key: str, value: str) -> bool: """释放锁(仅当持有者匹配时才删除)。""" r = await get_redis() script = """ if redis.call('get', KEYS[1]) == ARGV[1] then return redis.call('del', KEYS[1]) else return 0 end """ result = await r.eval(script, 1, key, value) return bool(result)