cache.py 6.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208
  1. """FAQ sec_ids 共享缓存(基于 Redis) + Leader 定时刷新。
  2. 设计要点:
  3. - 数据存 Redis 的 KEY_CACHE_DATA(JSON 格式),所有 worker 共享
  4. - 多 worker 启动时通过 KEY_LOCK_LEADER 选举出唯一 leader
  5. - 仅 leader 执行刷新 + 跑 scheduler_loop;follower 只读 Redis
  6. - leader 心跳续期;leader 进程挂掉后锁自然过期,其他 worker 抢锁接管
  7. """
  8. from __future__ import annotations
  9. import asyncio
  10. import json
  11. import logging
  12. import os
  13. import uuid
  14. from datetime import datetime, time, timedelta
  15. from typing import Any
  16. from app.config import settings
  17. from app.services.redis_client import (
  18. KEY_CACHE_DATA,
  19. KEY_LOCK_LEADER,
  20. KEY_LOCK_REFRESH,
  21. get_redis,
  22. release_lock,
  23. renew_lock,
  24. try_acquire_lock,
  25. )
  26. from app.services.zendesk_client import list_all_faq_sections
  27. logger = logging.getLogger(__name__)
  28. # 当前进程的 leader 标识(用于锁续期校验)
  29. WORKER_ID = f"{os.getpid()}-{uuid.uuid4().hex[:8]}"
  30. # 空快照(Redis 不可用或缓存未就绪时使用,保证业务不崩)
  31. _EMPTY_SNAPSHOT: dict[str, Any] = {
  32. "sec_ids": [],
  33. "faq_sections": [],
  34. "last_updated_at": None,
  35. "next_refresh_at": None,
  36. }
  37. def _next_refresh_time(now: datetime) -> datetime:
  38. """计算下一次定时刷新时间(每天 hh:mm,本地时区)。"""
  39. target_time = time(
  40. hour=settings.cache_refresh_hour,
  41. minute=settings.cache_refresh_minute,
  42. )
  43. today_target = datetime.combine(now.date(), target_time)
  44. if now < today_target:
  45. return today_target
  46. return today_target + timedelta(days=1)
  47. async def get_snapshot() -> dict[str, Any]:
  48. """从 Redis 读快照;Redis 不可用或无数据时返回空快照。"""
  49. try:
  50. r = await get_redis()
  51. raw = await r.get(KEY_CACHE_DATA)
  52. if not raw:
  53. return dict(_EMPTY_SNAPSHOT)
  54. return json.loads(raw)
  55. except Exception as exc: # noqa: BLE001
  56. # Redis 故障不应影响搜索请求,降级为空过滤
  57. logger.warning("读取 Redis 缓存失败,使用空快照: %s", exc)
  58. return dict(_EMPTY_SNAPSHOT)
  59. async def refresh_cache(force: bool = False) -> bool:
  60. """从 Zendesk 拉取最新数据并写入 Redis。
  61. 使用 KEY_LOCK_REFRESH 防并发刷新(避免重复打 Zendesk)。
  62. Args:
  63. force: True 时跳过去重锁直接刷新(管理员手动触发)
  64. Returns:
  65. True 实际执行了刷新;False 被去重锁挡住
  66. """
  67. if not force:
  68. got = await try_acquire_lock(
  69. KEY_LOCK_REFRESH, WORKER_ID, ttl_seconds=30
  70. )
  71. if not got:
  72. logger.info("另一进程正在刷新缓存,本次跳过")
  73. return False
  74. try:
  75. logger.info("[%s] 开始刷新 FAQ 缓存…", WORKER_ID)
  76. try:
  77. sec_ids, faq_sections = await list_all_faq_sections()
  78. except Exception as exc: # noqa: BLE001
  79. logger.exception("FAQ 缓存刷新失败:%s", exc)
  80. return False
  81. if not sec_ids:
  82. logger.warning("Zendesk 没有 FAQ sec_ids,无法刷新缓存")
  83. return False
  84. now = datetime.now()
  85. payload = {
  86. "sec_ids": list(sec_ids),
  87. "faq_sections": list(faq_sections),
  88. "last_updated_at": now.isoformat(),
  89. "next_refresh_at": _next_refresh_time(now).isoformat(),
  90. }
  91. try:
  92. r = await get_redis()
  93. await r.set(KEY_CACHE_DATA, json.dumps(payload, ensure_ascii=False))
  94. except Exception as exc: # noqa: BLE001
  95. logger.exception("写入 Redis 缓存失败: %s", exc)
  96. return False
  97. logger.info(
  98. "[%s] FAQ 缓存刷新完成:sec_ids=%s 条,下一次刷新 %s",
  99. WORKER_ID,
  100. len(payload["sec_ids"]),
  101. payload["next_refresh_at"],
  102. )
  103. return True
  104. finally:
  105. if not force:
  106. await release_lock(KEY_LOCK_REFRESH, WORKER_ID)
  107. # ===== Leader 选举与心跳 =====
  108. async def try_become_leader() -> bool:
  109. """尝试成为 leader(仅一个 worker 会成功)。"""
  110. try:
  111. return await try_acquire_lock(
  112. KEY_LOCK_LEADER, WORKER_ID, ttl_seconds=settings.leader_lock_ttl
  113. )
  114. except Exception as exc: # noqa: BLE001
  115. logger.warning("Leader 选举异常: %s", exc)
  116. return False
  117. async def leader_heartbeat() -> None:
  118. """Leader 心跳:周期性续期 leader 锁。
  119. 锁过期则尝试重新抢锁;抢到继续当 leader,抢不到说明已被他人接管,退出。
  120. """
  121. interval = max(1, settings.leader_lock_ttl // 3)
  122. while True:
  123. try:
  124. await asyncio.sleep(interval)
  125. except asyncio.CancelledError:
  126. raise
  127. try:
  128. ok = await renew_lock(
  129. KEY_LOCK_LEADER, WORKER_ID, settings.leader_lock_ttl
  130. )
  131. if not ok:
  132. # 锁丢失,尝试重新抢
  133. logger.warning("[%s] Leader 锁丢失,尝试重新选举", WORKER_ID)
  134. regained = await try_become_leader()
  135. if not regained:
  136. logger.info("[%s] 已不再是 leader,停止心跳", WORKER_ID)
  137. return
  138. logger.info("[%s] 重新成为 leader", WORKER_ID)
  139. except Exception as exc: # noqa: BLE001
  140. logger.warning("Leader 心跳异常: %s", exc)
  141. async def scheduler_loop() -> None:
  142. """Leader 专属:每天 hh:mm 刷新一次 FAQ 缓存。"""
  143. while True:
  144. now = datetime.now()
  145. next_run = _next_refresh_time(now)
  146. wait_seconds = (next_run - now).total_seconds()
  147. logger.info(
  148. "[%s] 下一次 FAQ 缓存刷新 %s(%.0f 秒后)",
  149. WORKER_ID,
  150. next_run.isoformat(),
  151. wait_seconds,
  152. )
  153. try:
  154. await asyncio.sleep(wait_seconds)
  155. except asyncio.CancelledError:
  156. logger.info("scheduler_loop 收到取消信号,退出")
  157. raise
  158. await refresh_cache()
  159. async def wait_cache_ready(timeout_seconds: int | None = None) -> bool:
  160. """Follower 启动时等待 leader 把缓存写入 Redis。
  161. Returns:
  162. True 缓存已就绪;False 超时仍未就绪(应用照样启动,搜索降级)
  163. """
  164. timeout = timeout_seconds or settings.cache_ready_wait
  165. deadline = asyncio.get_event_loop().time() + timeout
  166. while asyncio.get_event_loop().time() < deadline:
  167. try:
  168. r = await get_redis()
  169. exists = await r.exists(KEY_CACHE_DATA)
  170. if exists:
  171. return True
  172. except Exception as exc: # noqa: BLE001
  173. logger.warning("等待缓存就绪时 Redis 异常: %s", exc)
  174. return False
  175. await asyncio.sleep(0.5)
  176. return False