dedup.py 1.7 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950
  1. """message_id 幂等去重:基于 OrderedDict 实现 LRU + 容量上限。
  2. 飞书在网络抖动时会重发事件推送,不能依赖 event_id,必须用 message_id 本地去重。
  3. """
  4. from __future__ import annotations
  5. import logging
  6. import threading
  7. from collections import OrderedDict
  8. from typing import Optional
  9. logger = logging.getLogger(__name__)
  10. class MessageDedup:
  11. """线程安全的 LRU message_id 去重器。
  12. 飞书事件推送在网络抖动时会重发,必须用 message_id 在本地去重,
  13. 不能依赖 event_id(官方文档明确说明 event_id 不保证幂等)。
  14. """
  15. def __init__(self, max_size: int = 2000) -> None:
  16. if max_size <= 0:
  17. raise ValueError("max_size 必须为正数")
  18. self._max_size = max_size
  19. self._seen: OrderedDict[str, None] = OrderedDict()
  20. self._lock = threading.Lock()
  21. def check_and_mark(self, message_id: str) -> bool:
  22. """若 message_id 首次出现则记录并返回 True;已存在则返回 False(重复)。"""
  23. if not message_id:
  24. return False
  25. with self._lock:
  26. if message_id in self._seen:
  27. # 命中:移到末尾(LRU)
  28. self._seen.move_to_end(message_id)
  29. return False
  30. self._seen[message_id] = None
  31. if len(self._seen) > self._max_size:
  32. evicted_k, _ = self._seen.popitem(last=False)
  33. logger.debug("去重缓存淘汰最旧 message_id: %s", evicted_k)
  34. return True
  35. def __len__(self) -> int:
  36. with self._lock:
  37. return len(self._seen)
  38. def stats(self) -> dict:
  39. with self._lock:
  40. return {"size": len(self._seen), "max_size": self._max_size}