| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293 |
- """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.core.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)
|