feat: add global semaphore to limit concurrent snapshot builds

This commit is contained in:
HibiKier
2025-12-29 16:35:12 +08:00
parent cd2fd77789
commit 94939d4665
2 changed files with 77 additions and 20 deletions
+29 -6
View File
@@ -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 次
```
---
+48 -14
View File
@@ -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)