diff --git a/zhenxun/builtin_plugins/hooks/auth/auth_ban.py b/zhenxun/builtin_plugins/hooks/auth/auth_ban.py index 7eea7f57..0ab3cb70 100644 --- a/zhenxun/builtin_plugins/hooks/auth/auth_ban.py +++ b/zhenxun/builtin_plugins/hooks/auth/auth_ban.py @@ -79,21 +79,37 @@ async def is_ban(user_id: str | None, group_id: str | None) -> int: ban_dao.safe_get_or_none(user_id=user_id, group_id__isnull=True) ) - # 等待所有查询完成,添加超时控制 + # 等待所有查询完成,添加超时控制(使用更短的超时时间,快速失败) if tasks: try: + # 使用更短的超时时间(1.5秒),避免在高并发下等待太久 + # 如果查询超时,视为未ban,允许继续执行 ban_records = await asyncio.wait_for( - asyncio.gather(*tasks), timeout=DB_TIMEOUT_SECONDS + asyncio.gather(*tasks, return_exceptions=True), + timeout=min(DB_TIMEOUT_SECONDS, 1.5), ) + # 处理可能的异常 + valid_records = [] + for record in ban_records: + if isinstance(record, Exception): + logger.warning( + f"查询ban记录时出现异常: {record}", + LOGGER_COMMAND, + ) + continue + valid_records.append(record) + if len(tasks) == 2: - group_user, user = ban_records + group_user = valid_records[0] if len(valid_records) > 0 else None + user = valid_records[1] if len(valid_records) > 1 else None elif user_id and group_id: - group_user = ban_records[0] + group_user = valid_records[0] if valid_records else None else: - user = ban_records[0] + user = valid_records[0] if valid_records else None except asyncio.TimeoutError: - logger.error( - f"查询ban记录超时: user_id={user_id}, group_id={group_id}", + logger.warning( + f"查询ban记录超时(视为未ban): " + f"user_id={user_id}, group_id={group_id}", LOGGER_COMMAND, ) return 0 diff --git a/zhenxun/services/cache/__init__.py b/zhenxun/services/cache/__init__.py index 9e222a44..aba34539 100644 --- a/zhenxun/services/cache/__init__.py +++ b/zhenxun/services/cache/__init__.py @@ -599,8 +599,6 @@ class CacheManager: 返回: bool: 是否成功 """ - from zhenxun.services.db_context import DB_TIMEOUT_SECONDS - # 如果缓存被禁用或缓存模式为NONE,直接返回False if not self.enabled or cache_config.cache_mode == CacheMode.NONE: return False @@ -615,14 +613,17 @@ class CacheManager: # 设置过期时间 ttl = expire if expire is not None else model.expire - # 设置缓存 + # 设置缓存(使用较短的超时时间,避免阻塞主流程) await asyncio.wait_for( self.cache_backend.set(cache_key, serialized_value, ttl=ttl), # type: ignore - timeout=DB_TIMEOUT_SECONDS, + timeout=min(CACHE_TIMEOUT, 2.0), # 最多2秒,避免阻塞太久 ) return True except asyncio.TimeoutError: - logger.error(f"设置缓存 {cache_type}:{cache_key} 超时", LOG_COMMAND) + logger.warning( + f"设置缓存 {cache_type}:{cache_key} 超时(已跳过,不影响主流程)", + LOG_COMMAND, + ) return False except Exception as e: logger.error(f"设置缓存 {cache_type} 失败", LOG_COMMAND, e=e) diff --git a/zhenxun/services/data_access.py b/zhenxun/services/data_access.py index 3912fe28..a3d0216a 100644 --- a/zhenxun/services/data_access.py +++ b/zhenxun/services/data_access.py @@ -1,3 +1,4 @@ +import asyncio from typing import Any, ClassVar, Generic, TypeVar, cast from zhenxun.services.cache import Cache, CacheRoot, cache_config @@ -212,9 +213,13 @@ class DataAccess(Generic[T]): except Exception as e: logger.error(f"{self.model_cls.__name__} 从缓存获取数据失败: {kwargs}", e=e) - # 如果缓存中没有,从数据库获取 + # 如果缓存中没有,从数据库获取(使用超时控制) logger.debug(f"{self.model_cls.__name__} 从数据库获取数据: {kwargs}") - data = await db_query_func(*args, **kwargs) + data = await with_db_timeout( + db_query_func(*args, **kwargs), + operation=f"{self.model_cls.__name__}.{db_query_func.__name__}", + source="DataAccess._get_with_cache", + ) # 如果获取到数据,存入缓存 if data: @@ -222,31 +227,48 @@ class DataAccess(Generic[T]): # 生成缓存键 cache_key = self._build_cache_key_for_item(data) if cache_key is not None: - # 存入缓存 - await self.cache.set(cache_key, data) - self._cache_stats[self.cache_type]["sets"] += 1 - logger.debug( - f"{self.model_cls.__name__} 数据已存入缓存: {cache_key}" - ) + # 存入缓存(失败不影响主流程) + try: + # 使用较短的超时时间,避免阻塞 + await asyncio.wait_for( + self.cache.set(cache_key, data), timeout=1.0 + ) + self._cache_stats[self.cache_type]["sets"] += 1 + logger.debug( + f"{self.model_cls.__name__} 数据已存入缓存: {cache_key}" + ) + except (asyncio.TimeoutError, Exception) as cache_err: + # 缓存设置失败不影响数据返回,只记录警告 + logger.warning( + f"{self.model_cls.__name__} 存入缓存失败(超时或异常)," + f"参数: {kwargs}", + e=cache_err, + ) except Exception as e: logger.error( f"{self.model_cls.__name__} 存入缓存失败,参数: {kwargs}", e=e ) elif cache_key is not None: - # 如果没有获取到数据,缓存空结果 + # 如果没有获取到数据,缓存空结果(失败不影响主流程) try: - # 存入空结果缓存,使用较短的过期时间 - await self.cache.set( - cache_key, self._NULL_RESULT, expire=self._NULL_RESULT_TTL + # 存入空结果缓存,使用较短的过期时间和超时时间 + await asyncio.wait_for( + self.cache.set( + cache_key, self._NULL_RESULT, expire=self._NULL_RESULT_TTL + ), + timeout=1.0, ) self._cache_stats[self.cache_type]["null_sets"] += 1 logger.debug( f"{self.model_cls.__name__} 空结果已存入缓存: {cache_key}," f" TTL={self._NULL_RESULT_TTL}秒" ) - except Exception as e: - logger.error( - f"{self.model_cls.__name__} 存入空结果缓存失败,参数: {kwargs}", e=e + except (asyncio.TimeoutError, Exception) as cache_err: + # 空结果缓存设置失败不影响数据返回,只记录警告 + logger.warning( + f"{self.model_cls.__name__} 存入空结果缓存失败(超时或异常)," + f"参数: {kwargs}", + e=cache_err, ) return data