| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208 |
- """FAQ sec_ids 共享缓存(基于 Redis) + Leader 定时刷新。
- 设计要点:
- - 数据存 Redis 的 KEY_CACHE_DATA(JSON 格式),所有 worker 共享
- - 多 worker 启动时通过 KEY_LOCK_LEADER 选举出唯一 leader
- - 仅 leader 执行刷新 + 跑 scheduler_loop;follower 只读 Redis
- - leader 心跳续期;leader 进程挂掉后锁自然过期,其他 worker 抢锁接管
- """
- from __future__ import annotations
- import asyncio
- import json
- import logging
- import os
- import uuid
- from datetime import datetime, time, timedelta
- from typing import Any
- from app.config import settings
- from app.services.redis_client import (
- KEY_CACHE_DATA,
- KEY_LOCK_LEADER,
- KEY_LOCK_REFRESH,
- get_redis,
- release_lock,
- renew_lock,
- try_acquire_lock,
- )
- from app.services.zendesk_client import list_all_faq_sections
- logger = logging.getLogger(__name__)
- # 当前进程的 leader 标识(用于锁续期校验)
- WORKER_ID = f"{os.getpid()}-{uuid.uuid4().hex[:8]}"
- # 空快照(Redis 不可用或缓存未就绪时使用,保证业务不崩)
- _EMPTY_SNAPSHOT: dict[str, Any] = {
- "sec_ids": [],
- "faq_sections": [],
- "last_updated_at": None,
- "next_refresh_at": None,
- }
- def _next_refresh_time(now: datetime) -> datetime:
- """计算下一次定时刷新时间(每天 hh:mm,本地时区)。"""
- target_time = time(
- hour=settings.cache_refresh_hour,
- minute=settings.cache_refresh_minute,
- )
- today_target = datetime.combine(now.date(), target_time)
- if now < today_target:
- return today_target
- return today_target + timedelta(days=1)
- async def get_snapshot() -> dict[str, Any]:
- """从 Redis 读快照;Redis 不可用或无数据时返回空快照。"""
- try:
- r = await get_redis()
- raw = await r.get(KEY_CACHE_DATA)
- if not raw:
- return dict(_EMPTY_SNAPSHOT)
- return json.loads(raw)
- except Exception as exc: # noqa: BLE001
- # Redis 故障不应影响搜索请求,降级为空过滤
- logger.warning("读取 Redis 缓存失败,使用空快照: %s", exc)
- return dict(_EMPTY_SNAPSHOT)
- async def refresh_cache(force: bool = False) -> bool:
- """从 Zendesk 拉取最新数据并写入 Redis。
- 使用 KEY_LOCK_REFRESH 防并发刷新(避免重复打 Zendesk)。
- Args:
- force: True 时跳过去重锁直接刷新(管理员手动触发)
- Returns:
- True 实际执行了刷新;False 被去重锁挡住
- """
- if not force:
- got = await try_acquire_lock(
- KEY_LOCK_REFRESH, WORKER_ID, ttl_seconds=30
- )
- if not got:
- logger.info("另一进程正在刷新缓存,本次跳过")
- return False
- try:
- logger.info("[%s] 开始刷新 FAQ 缓存…", WORKER_ID)
- try:
- sec_ids, faq_sections = await list_all_faq_sections()
- except Exception as exc: # noqa: BLE001
- logger.exception("FAQ 缓存刷新失败:%s", exc)
- return False
- if not sec_ids:
- logger.warning("Zendesk 没有 FAQ sec_ids,无法刷新缓存")
- return False
- now = datetime.now()
- payload = {
- "sec_ids": list(sec_ids),
- "faq_sections": list(faq_sections),
- "last_updated_at": now.isoformat(),
- "next_refresh_at": _next_refresh_time(now).isoformat(),
- }
- try:
- r = await get_redis()
- await r.set(KEY_CACHE_DATA, json.dumps(payload, ensure_ascii=False))
- except Exception as exc: # noqa: BLE001
- logger.exception("写入 Redis 缓存失败: %s", exc)
- return False
- logger.info(
- "[%s] FAQ 缓存刷新完成:sec_ids=%s 条,下一次刷新 %s",
- WORKER_ID,
- len(payload["sec_ids"]),
- payload["next_refresh_at"],
- )
- return True
- finally:
- if not force:
- await release_lock(KEY_LOCK_REFRESH, WORKER_ID)
- # ===== Leader 选举与心跳 =====
- async def try_become_leader() -> bool:
- """尝试成为 leader(仅一个 worker 会成功)。"""
- try:
- return await try_acquire_lock(
- KEY_LOCK_LEADER, WORKER_ID, ttl_seconds=settings.leader_lock_ttl
- )
- except Exception as exc: # noqa: BLE001
- logger.warning("Leader 选举异常: %s", exc)
- return False
- async def leader_heartbeat() -> None:
- """Leader 心跳:周期性续期 leader 锁。
- 锁过期则尝试重新抢锁;抢到继续当 leader,抢不到说明已被他人接管,退出。
- """
- interval = max(1, settings.leader_lock_ttl // 3)
- while True:
- try:
- await asyncio.sleep(interval)
- except asyncio.CancelledError:
- raise
- try:
- ok = await renew_lock(
- KEY_LOCK_LEADER, WORKER_ID, settings.leader_lock_ttl
- )
- if not ok:
- # 锁丢失,尝试重新抢
- logger.warning("[%s] Leader 锁丢失,尝试重新选举", WORKER_ID)
- regained = await try_become_leader()
- if not regained:
- logger.info("[%s] 已不再是 leader,停止心跳", WORKER_ID)
- return
- logger.info("[%s] 重新成为 leader", WORKER_ID)
- except Exception as exc: # noqa: BLE001
- logger.warning("Leader 心跳异常: %s", exc)
- async def scheduler_loop() -> None:
- """Leader 专属:每天 hh:mm 刷新一次 FAQ 缓存。"""
- while True:
- now = datetime.now()
- next_run = _next_refresh_time(now)
- wait_seconds = (next_run - now).total_seconds()
- logger.info(
- "[%s] 下一次 FAQ 缓存刷新 %s(%.0f 秒后)",
- WORKER_ID,
- next_run.isoformat(),
- wait_seconds,
- )
- try:
- await asyncio.sleep(wait_seconds)
- except asyncio.CancelledError:
- logger.info("scheduler_loop 收到取消信号,退出")
- raise
- await refresh_cache()
- async def wait_cache_ready(timeout_seconds: int | None = None) -> bool:
- """Follower 启动时等待 leader 把缓存写入 Redis。
- Returns:
- True 缓存已就绪;False 超时仍未就绪(应用照样启动,搜索降级)
- """
- timeout = timeout_seconds or settings.cache_ready_wait
- deadline = asyncio.get_event_loop().time() + timeout
- while asyncio.get_event_loop().time() < deadline:
- try:
- r = await get_redis()
- exists = await r.exists(KEY_CACHE_DATA)
- if exists:
- return True
- except Exception as exc: # noqa: BLE001
- logger.warning("等待缓存就绪时 Redis 异常: %s", exc)
- return False
- await asyncio.sleep(0.5)
- return False
|