From 94939d466582c42d4fee792661238300d3d2f1be Mon Sep 17 00:00:00 2001 From: HibiKier <775757368@qq.com> Date: Mon, 29 Dec 2025 16:35:12 +0800 Subject: [PATCH] feat: add global semaphore to limit concurrent snapshot builds --- plan.md | 35 ++++++++++--- zhenxun/services/auth_snapshot/service.py | 62 ++++++++++++++++++----- 2 files changed, 77 insertions(+), 20 deletions(-) diff --git a/plan.md b/plan.md index b74a4eb5..d591c4a9 100644 --- a/plan.md +++ b/plan.md @@ -233,12 +233,35 @@ zhenxun/ ## 风险与缓解 -| 风险 | 缓解措施 | -| ------------ | ------------------------------- | -| 快照数据过期 | 合理的 TTL + 主动失效机制 | -| 快照构建延迟 | 异步构建 + 首次访问降级到旧流程 | -| 内存占用增加 | 监控内存使用 + 合理的缓存清理 | -| 数据一致性 | 写操作后立即失效缓存 | +| 风险 | 缓解措施 | +| ----------------- | -------------------------------------------- | +| 快照数据过期 | 合理的 TTL + 主动失效机制 | +| 快照构建延迟 | 异步构建 + 首次访问降级到旧流程 | +| 内存占用增加 | 监控内存使用 + 合理的缓存清理 | +| 数据一致性 | 写操作后立即失效缓存 | +| **DB 过载风险** | **全局 Semaphore 限制并发构建数量 (50)** | +| **并发构建重复** | **按 cache_key 的 asyncio.Lock** | +| **构建等待超时** | **3 秒超时后返回默认快照,允许请求继续处理** | + +### 并发控制机制 + +``` +大量消息同时进入时: + +1. 同一 user:group:bot 组合 + - 使用 asyncio.Lock 保证只构建一次 + - 其他等待的协程复用同一个 Future 结果 + +2. 不同 user:group:bot 组合 + - 使用全局 Semaphore 限制最多 50 个并发构建 + - 超过限制的请求排队等待(最多 3 秒) + - 等待超时则返回默认快照,避免请求阻塞 + +这样即使 1000 个不同用户同时发消息: +- 最多只有 50 个并发 DB 查询 +- 每个构建 5-7 次查询 = 最多 350 次并发 DB 查询 +- 远低于直接查询的 6000 次 +``` --- diff --git a/zhenxun/services/auth_snapshot/service.py b/zhenxun/services/auth_snapshot/service.py index 297b69a8..b8db6ffe 100644 --- a/zhenxun/services/auth_snapshot/service.py +++ b/zhenxun/services/auth_snapshot/service.py @@ -28,6 +28,10 @@ AUTH_REDIS_TTL = 60 # 权限快照Redis缓存TTL(秒) PLUGIN_MEMORY_TTL = 30 # 插件快照内存缓存TTL(秒) PLUGIN_REDIS_TTL = 300 # 插件快照Redis缓存TTL(秒) +# 并发控制配置 +MAX_CONCURRENT_BUILDS = 50 # 最大同时构建数量(防止 DB 过载) +BUILD_QUEUE_TIMEOUT = 3.0 # 等待构建队列的超时时间(秒) + class AuthSnapshotService: """权限快照服务 @@ -44,6 +48,16 @@ class AuthSnapshotService: # 构建锁(按 cache_key 粒度) _build_locks: ClassVar[dict[str, asyncio.Lock]] = {} + # 全局构建并发限制(防止大量不同 key 同时构建导致 DB 过载) + _build_semaphore: ClassVar[asyncio.Semaphore | None] = None + + @classmethod + def _get_build_semaphore(cls) -> asyncio.Semaphore: + """获取构建信号量(懒加载)""" + if cls._build_semaphore is None: + cls._build_semaphore = asyncio.Semaphore(MAX_CONCURRENT_BUILDS) + return cls._build_semaphore + @classmethod def _get_memory_cache(cls) -> CacheDict[AuthSnapshot]: """获取内存缓存实例(懒加载)""" @@ -133,28 +147,48 @@ class AuthSnapshotService: bot_id: str, cache_key: str, ) -> AuthSnapshot: - """构建并缓存快照""" + """构建并缓存快照(带全局并发限制)""" loop = asyncio.get_running_loop() future: asyncio.Future[AuthSnapshot] = loop.create_future() cls._building[cache_key] = future - try: - # 构建快照 - snapshot = await SnapshotBuilder.build_auth_snapshot( - user_id, group_id, bot_id - ) + semaphore = cls._get_build_semaphore() - # 存入Redis缓存(异步,不阻塞) - if cache_config.cache_mode != CacheMode.NONE: - asyncio.create_task( # noqa: RUF006 - cls._cache_to_redis(cache_key, snapshot) + try: + # 尝试获取信号量(限制并发构建数量) + try: + await asyncio.wait_for(semaphore.acquire(), timeout=BUILD_QUEUE_TIMEOUT) + except asyncio.TimeoutError: + # 等待超时,返回默认快照(允许请求继续,但不保证权限数据完整) + logger.warning( + f"构建快照等待超时(并发过高),使用默认快照: {cache_key}", + LOG_COMMAND, + ) + snapshot = AuthSnapshot( + user_id=user_id, group_id=group_id, bot_id=bot_id + ) + future.set_result(snapshot) + return snapshot + + try: + # 构建快照 + snapshot = await SnapshotBuilder.build_auth_snapshot( + user_id, group_id, bot_id ) - # 存入内存缓存 - cls._get_memory_cache().set(cache_key, snapshot) + # 存入Redis缓存(异步,不阻塞) + if cache_config.cache_mode != CacheMode.NONE: + asyncio.create_task( # noqa: RUF006 + cls._cache_to_redis(cache_key, snapshot) + ) - future.set_result(snapshot) - return snapshot + # 存入内存缓存 + cls._get_memory_cache().set(cache_key, snapshot) + + future.set_result(snapshot) + return snapshot + finally: + semaphore.release() except Exception as e: future.set_exception(e)