From faffdcad4708947ce15e46fb9e9c8e743012418a Mon Sep 17 00:00:00 2001 From: Kris <2893855659@qq.com> Date: Wed, 13 May 2026 16:19:39 +0800 Subject: [PATCH] =?UTF-8?q?feat(dedup=5Fpolicy):=20=E4=B8=BA=E5=8E=BB?= =?UTF-8?q?=E9=87=8D=E7=AD=96=E7=95=A5=E6=B7=BB=E5=8A=A0Redis=E5=90=8E?= =?UTF-8?q?=E7=AB=AF=E6=94=AF=E6=8C=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 重构了去重策略的核心逻辑,新增Redis分布式去重能力,当配置redis_url时优先使用Redis存储去重状态, fallback到内存缓存;同时调整了方法顺序并新增close方法用于资源清理 --- .../yuxi/channels/policy/dedup_policy.py | 62 ++++++++++++++++--- 1 file changed, 55 insertions(+), 7 deletions(-) diff --git a/backend/package/yuxi/channels/policy/dedup_policy.py b/backend/package/yuxi/channels/policy/dedup_policy.py index aeff4354..fde01d0a 100644 --- a/backend/package/yuxi/channels/policy/dedup_policy.py +++ b/backend/package/yuxi/channels/policy/dedup_policy.py @@ -8,21 +8,61 @@ from yuxi.utils.logging_config import logger class DedupPolicy: - def __init__(self, ttl: int = 300, maxsize: int = 10000): + def __init__(self, ttl: int = 300, maxsize: int = 10000, redis_url: str | None = None): self._seen: TTLCache = TTLCache(maxsize=maxsize, ttl=ttl) self._lock = asyncio.Lock() + self._redis_url = redis_url + self._redis = None + + async def _ensure_redis(self): + if self._redis is not None: + return self._redis + if not self._redis_url: + return None + try: + import redis.asyncio as aioredis + + self._redis = aioredis.from_url(self._redis_url, decode_responses=False) + await self._redis.ping() + logger.info(f"DedupPolicy: Redis connected ({self._redis_url})") + return self._redis + except Exception as e: + logger.warning(f"DedupPolicy: Redis unavailable, falling back to memory: {e}") + if self._redis is not None: + try: + await self._redis.aclose() + except Exception: + pass + self._redis = None + return None + + async def check_and_mark(self, key: str, ttl: int | None = None) -> bool: + effective_ttl = ttl if ttl is not None else self._seen.ttl + + redis = await self._ensure_redis() + if redis is not None: + try: + acquired = await redis.set(key, "1", nx=True, ex=effective_ttl) + if acquired is None: + return True + async with self._lock: + self._seen[key] = True + return False + except Exception as e: + logger.debug(f"DedupPolicy: Redis error, falling back to memory: {e}") + + async with self._lock: + if key in self._seen: + return True + self._seen[key] = True + return False async def is_duplicate(self, message: ChannelMessage) -> bool: msg_id = message.identity.channel_message_id if not msg_id: return False key = f"{message.identity.channel_id}:{msg_id}" - async with self._lock: - if key in self._seen: - logger.debug(f"Duplicate message filtered: {key}") - return True - self._seen[key] = True - return False + return await self.check_and_mark(key) async def check_and_remember(self, message: ChannelMessage) -> bool: async with self._lock: @@ -77,3 +117,11 @@ class DedupPolicy: async def clear(self) -> None: async with self._lock: self._seen.clear() + + async def close(self) -> None: + if self._redis is not None: + try: + await self._redis.aclose() + except Exception: + pass + self._redis = None