Răsfoiți Sursa

优化多进程启动逻辑

liujintao 1 lună în urmă
părinte
comite
79834053ef
12 a modificat fișierele cu 644 adăugiri și 141 ștergeri
  1. 6 0
      .env
  2. 27 0
      .env.dev.example
  3. 16 0
      .env.example
  4. 8 0
      .gitignore
  5. 198 45
      README.md
  6. 33 2
      app/config.py
  7. 58 18
      app/main.py
  8. 3 3
      app/routers/search.py
  9. 166 64
      app/services/cache.py
  10. 93 0
      app/services/redis_client.py
  11. 35 9
      app/services/zendesk_client.py
  12. 1 0
      requirements.txt

+ 6 - 0
.env

@@ -3,6 +3,12 @@ ZENDESK_SUBDOMAIN=zositechhelp
 ZENDESK_EMAIL=zhzxb027@vip.163.com
 ZENDESK_API_TOKEN=CpmJL8D4aP20BPn9vi8t9ww5VUOv0bhhu69PFkyc
 
+REDIS_URL=redis://retail.zafjfx.ng.0001.use1.cache.amazonaws.com:6379/0
+REDIS_PASSWORD=
+REDIS_KEY_PREFIX=faq_search
+LEADER_LOCK_TTL=60
+CACHE_READY_WAIT=180
+
 # 默认语言
 DEFAULT_LOCALE=en-us
 

+ 27 - 0
.env.dev.example

@@ -0,0 +1,27 @@
+# ============================================
+# 开发环境配置(APP_ENV=dev 时加载此文件)
+# 此文件已被 .gitignore 忽略,请填入你的本地开发凭据
+# ============================================
+
+# Zendesk Help Center 凭据(开发环境建议使用 sandbox / 测试账号)
+ZENDESK_SUBDOMAIN=zositechhelp
+ZENDESK_EMAIL=your_dev_email@example.com
+ZENDESK_API_TOKEN=your_dev_api_token
+
+# Redis(本地开发用 docker run -p 6379:6379 redis 起一个就行)
+# 注意:开发环境用 db=1 与生产环境隔离,避免污染生产缓存
+REDIS_URL=redis://127.0.0.1:6379/1
+REDIS_PASSWORD=
+REDIS_KEY_PREFIX=faq_search_dev
+LEADER_LOCK_TTL=60
+CACHE_READY_WAIT=30
+
+# 默认语言
+DEFAULT_LOCALE=en-us
+
+# FAQ 缓存刷新时间(开发环境一般用 /admin/cache/refresh 手动触发,定时设大点)
+CACHE_REFRESH_HOUR=0
+CACHE_REFRESH_MINUTE=0
+
+# HTTP 请求超时(秒)
+HTTP_TIMEOUT=30

+ 16 - 0
.env.example

@@ -1,8 +1,23 @@
+# ============================================
+# 生产环境配置模板(APP_ENV=prod 或不设置时加载 .env)
+# 复制为 .env 并填入真实凭据;.env 已被 .gitignore 忽略
+# 开发环境请复制 .env.dev.example -> .env.dev 并使用 APP_ENV=dev 启动
+# ============================================
+
 # Zendesk Help Center 凭据
 ZENDESK_SUBDOMAIN=zositechhelp
 ZENDESK_EMAIL=your_email@example.com
 ZENDESK_API_TOKEN=your_api_token_here
 
+# Redis(多 worker 共享缓存与 leader 选举,必填)
+REDIS_URL=redis://127.0.0.1:6379/0
+REDIS_PASSWORD=
+REDIS_KEY_PREFIX=faq_search
+# leader 锁 TTL(秒),心跳每 1/3 续期一次;leader 挂掉后 TTL 秒内被替换
+LEADER_LOCK_TTL=60
+# follower 等待 leader 首次写入缓存的最大秒数;超时降级为空 sec_ids
+CACHE_READY_WAIT=180
+
 # 默认语言
 DEFAULT_LOCALE=en-us
 
@@ -12,3 +27,4 @@ CACHE_REFRESH_MINUTE=0
 
 # HTTP 请求超时(秒)
 HTTP_TIMEOUT=30
+

+ 8 - 0
.gitignore

@@ -15,3 +15,11 @@ dist
 .webpack-cache
 .serverless_plugins
 .venv
+
+# 环境变量(不进版本库)
+.env.dev
+.env.prod
+.env.local
+__pycache__/
+.mypy_cache/
+.pytest_cache/

+ 198 - 45
README.md

@@ -1,50 +1,104 @@
 # Zendesk FAQ Keyword Search
 
-把 `test1.py` (search_articles) 与 `test2.py` (list_all_sections_full) 重写为 FastAPI 服务:
+基于 FastAPI 的 Zendesk Help Center FAQ 搜索代理,封装 Zendesk Help Center API,自动维护 FAQ Section 缓存并提供统一的关键词搜索接口。
 
-- 启动时加载所有 **FAQ Section IDs**(名字含 `faq` 且文章数 > 0)保存到内存
-- 每天 **0:00** 自动刷新一次(可在 `.env` 改时间)
-- 暴露 `GET/POST /search` 接口,请求字段为 `query`,返回 Zendesk 接口的搜索结果
+## 特性
+
+- **FAQ Section 自动发现**:扫描 Zendesk Help Center 找出名字含 `faq` 且文章数 > 0 的 Section
+- **Redis 共享缓存**:多 worker 进程共享缓存,避免重复打 Zendesk
+- **Leader 选举**:仅一个 worker 负责刷新缓存,其余 worker 共享读取
+- **每日定时刷新**:可配置时间(默认 0:00)自动刷新 sec_ids
+- **故障自愈**:leader 进程崩溃后,其他 worker 自动接管
+- **统一响应格式**:所有接口返回 `{code, data, msg}` 信封
+- **优雅降级**:Redis 不可用时返回空 sec_ids,搜索不带 section 过滤继续工作
 
 ## 目录结构
 
 ```
 app/
-├── main.py               # FastAPI 入口、lifespan、健康检查、手动刷新
-├── config.py             # 从 .env 读取 Zendesk 凭据
-├── schemas.py            # 请求/响应 Pydantic 模型
-├── services/
-│   ├── zendesk_client.py # 异步 list_all_faq_sections / search_articles
-│   └── cache.py          # 内存缓存 + 每日 0 点定时刷新
-└── routers/
-    └── search.py         # /search 接口
+├── main.py                    # FastAPI 入口、lifespan、leader 选举、健康检查
+├── config.py                  # 配置加载(按 APP_ENV 切换 .env / .env.dev)
+├── schemas.py                 # 请求/响应 Pydantic 模型
+├── response.py                # 统一响应类 + 业务错误码 + BusinessError
+├── exception_handlers.py      # 全局异常处理器
+├── routers/
+│   └── search.py              # /search 接口
+└── services/
+    ├── zendesk_client.py      # 异步 Zendesk 客户端
+    ├── redis_client.py        # Redis 连接 + 分布式锁原语
+    └── cache.py               # 共享缓存读写 + leader 选举/心跳/定时刷新
 ```
 
-## 安装
+## 快速开始
+
+### 1. 安装依赖
 
 ```bash
 pip install -r requirements.txt
+```
+
+### 2. 准备 Redis
+
+```bash
+# Docker 一行起
+docker run -d -p 6379:6379 --name redis redis:7
+# 或用 brew (macOS)
+brew install redis && brew services start redis
+```
+
+### 3. 配置环境
+
+#### 生产环境
+
+```bash
 cp .env.example .env
-# 编辑 .env 填入 ZENDESK_SUBDOMAIN / ZENDESK_EMAIL / ZENDESK_API_TOKEN
+# 编辑 .env 填入真实 Zendesk 凭据和 Redis 密码
 ```
 
-## 运行
+#### 开发环境
 
 ```bash
-uvicorn app.main:app --host 0.0.0.0 --port 8000 --reload
+cp .env.dev.example .env.dev
+# 编辑 .env.dev 填入开发凭据
+# 注意:开发环境用 db=1 与生产隔离(REDIS_URL=redis://127.0.0.1:6379/1)
 ```
 
-启动日志中会看到:
+### 4. 启动
 
+```bash
+# 开发环境(自动重载、单 worker)
+APP_ENV=dev uvicorn app.main:app --host 127.0.0.1 --port 8000 --reload
+
+# 生产环境(多 worker)
+APP_ENV=prod uvicorn app.main:app --host 0.0.0.0 --port 8025 --workers 3
+# 或省略 APP_ENV(默认就是 prod)
+uvicorn app.main:app --host 0.0.0.0 --port 8025 --workers 3
 ```
-启动:首次加载 FAQ 缓存…
-✅ FAQ Section id=... count=... name=... cat=...
-FAQ 缓存刷新完成:sec_ids=N 条,下一次刷新 2026-06-16T00:00:00
-启动:定时刷新任务已运行
+
+启动日志中会看到 leader 选举:
+
+```
+[14408-abc] 当选 leader,开始首次刷新
+[14408-abc] 开始刷新 FAQ 缓存…
+[14409-def] 不是 leader,等待缓存就绪…
+[14410-ghi] 不是 leader,等待缓存就绪…
+[14408-abc] FAQ 缓存刷新完成:sec_ids=N 条,下一次刷新 ...
+[14409-def] 共享缓存已就绪
+[14410-ghi] 共享缓存已就绪
 ```
 
 ## 接口
 
+所有接口返回统一格式:
+
+```json
+{
+  "code": 0,
+  "data": {},
+  "msg": "成功"
+}
+```
+
 ### 1. 搜索(GET)
 
 ```bash
@@ -59,56 +113,155 @@ curl -X POST http://localhost:8000/search \
   -d '{"query": "Can I Add a PoE Switch", "page": 1, "per_page": 25}'
 ```
 
-响应:
+成功响应:
 
 ```json
 {
-  "success": true,
-  "query": "...",
-  "count": 3,
-  "page": 1,
-  "per_page": 25,
-  "next_page": null,
-  "sec_ids_used": [48252627506585, 47684283137817, "..."],
-  "results": [{"title": "...", "html_url": "...", "snippet": "..."}]
+  "code": 0,
+  "data": {
+    "query": "Can I Add a PoE Switch",
+    "count": 3,
+    "page": 1,
+    "per_page": 25,
+    "next_page": null,
+    "sec_ids_used": [48252627506585, 47684283137817],
+    "results": [{"title": "...", "html_url": "...", "snippet": "..."}]
+  },
+  "msg": "成功"
 }
 ```
 
-### 3. 缓存状态
+### 3. 健康检查
+
+```bash
+curl http://localhost:8000/health
+# {"code": 0, "data": {"status": "ok"}, "msg": "成功"}
+```
+
+### 4. 缓存状态
 
 ```bash
 curl http://localhost:8000/health/cache
-# {"sec_ids_count": 30, "last_updated_at": "2026-06-15T09:12:33", "next_refresh_at": "2026-06-16T00:00:00"}
 ```
 
-### 4. 手动刷新缓存
+```json
+{
+  "code": 0,
+  "data": {
+    "sec_ids_count": 30,
+    "last_updated_at": "2026-06-15T09:12:33",
+    "next_refresh_at": "2026-06-16T00:00:00"
+  },
+  "msg": "成功"
+}
+```
+
+### 5. 手动刷新缓存
 
 ```bash
 curl -X POST http://localhost:8000/admin/cache/refresh
 ```
 
-### 5. Swagger UI
+### 6. Swagger UI
 
 打开 <http://localhost:8000/docs>。
 
-## 配置项(.env)
+## 错误码表
+
+| code | 含义 | HTTP 状态 |
+|------|------|-----------|
+| 0 | 成功 | 200 |
+| 1001 | 参数验证失败(如 page=-1) | 422 |
+| 1002 | 请求格式不正确 | 400 |
+| 1003 | 请求方法不允许 | 405 |
+| 1004 | 资源不存在 | 404 |
+| 2001 | FAQ 缓存未就绪 | 200 |
+| 2002 | 业务资源不存在 | 200 |
+| 3001 | 上游服务调用失败(Zendesk) | 200 |
+| 3002 | 上游服务响应超时 | 200 |
+| 9999 | 服务器内部错误 | 500 |
+
+参数错误响应示例(`?query=hi&page=-1`):
+
+```json
+{
+  "code": 1001,
+  "data": [{"loc": ["query", "page"], "msg": "Input should be greater than or equal to 1", "type": "greater_than_equal"}],
+  "msg": "参数验证失败"
+}
+```
+
+## 配置项
+
+通过 `APP_ENV` 选择加载哪个 env 文件:
+
+| `APP_ENV` | 加载文件 |
+|-----------|----------|
+| `dev` | `.env.dev` |
+| `prod`(默认) | `.env` |
+
+### 完整配置项
 
 | 变量 | 默认 | 说明 |
 | --- | --- | --- |
 | `ZENDESK_SUBDOMAIN` | — | Zendesk 子域名,例:`zositechhelp` |
 | `ZENDESK_EMAIL` | — | Zendesk 账号邮箱 |
-| `ZENDESK_API_TOKEN` | — | API Token |
+| `ZENDESK_API_TOKEN` | — | Zendesk API Token |
 | `DEFAULT_LOCALE` | `en-us` | 默认语言 |
-| `CACHE_REFRESH_HOUR` | `0` | 每日刷新小时(本地时区) |
+| `CACHE_REFRESH_HOUR` | `0` | 每日刷新小时(24 小时制,本地时区) |
 | `CACHE_REFRESH_MINUTE` | `0` | 每日刷新分钟 |
 | `HTTP_TIMEOUT` | `30` | httpx 超时秒数 |
+| `REDIS_URL` | `redis://127.0.0.1:6379/0` | Redis 连接 URL |
+| `REDIS_PASSWORD` | — | Redis 密码(如有) |
+| `REDIS_USERNAME` | — | Redis 用户名(ACL 模式才需要) |
+| `REDIS_KEY_PREFIX` | `faq_search` | 所有 Redis Key 的前缀 |
+| `LEADER_LOCK_TTL` | `60` | leader 锁 TTL,决定故障切换延迟 |
+| `CACHE_READY_WAIT` | `30` | follower 等待缓存就绪的最大秒数 |
+
+### 环境隔离建议
+
+| 项 | 生产 | 开发 |
+|----|------|------|
+| Redis db | `0` | `1`(避免污染生产数据) |
+| `REDIS_KEY_PREFIX` | `faq_search` | `faq_search_dev` |
+| 启动 | `--workers 3` | `--reload`(单 worker) |
+| Zendesk 凭据 | 生产账号 | sandbox / 测试账号 |
 
-## 与原脚本对照
+## 多 Worker 与 Leader 机制
 
-| 原脚本 | 新位置 |
-| --- | --- |
-| `test2.list_all_sections_full` | `app/services/zendesk_client.list_all_faq_sections`(异步) |
-| `test2.get_article_count` | `app/services/zendesk_client._get_article_count`(私有) |
-| `test1.search_articles` | `app/services/zendesk_client.search_articles`(异步) |
-| `test1.faq_section_ids` 硬编码 | 内存缓存 `app/services/cache.faq_cache.sec_ids`(启动加载 + 每日 0 点刷新) |
-| `test1.QUERY` 硬编码 | HTTP 请求字段 `query` |
+启动 N 个 worker 时(`uvicorn --workers N`):
+
+1. 每个 worker 独立尝试在 Redis 抢 `lock:leader`,**仅一个**抢到
+2. Leader:拉 Zendesk 写入 Redis、跑 `scheduler_loop`(每天定时刷新)、跑 `leader_heartbeat`(每 `LEADER_LOCK_TTL/3` 秒续期锁)
+3. Follower:仅读 Redis 缓存,不打 Zendesk
+4. Leader 进程崩溃后,`LEADER_LOCK_TTL` 秒内剩余 worker 接管
+
+**Redis Key**:
+
+```
+faq_search:cache:data       JSON 缓存数据(共享)
+faq_search:lock:leader      leader 选举锁
+faq_search:lock:refresh     防并发刷新锁
+```
+
+## Redis 故障降级
+
+- **Redis 完全挂掉**:`get_snapshot()` 返回空快照,搜索不带 section 过滤继续工作(结果可能包含非 FAQ 文章)
+- **Leader 长时间卡死**:`LEADER_LOCK_TTL` 后锁过期,其他 worker 抢锁接管
+- **手动刷新**:`POST /admin/cache/refresh` 用 `force=True` 跳过去重锁,立即生效
+
+## 开发常见命令
+
+```bash
+# 仅启动开发环境
+APP_ENV=dev uvicorn app.main:app --reload
+
+# 用 redis-cli 看缓存
+redis-cli -n 1 get faq_search_dev:cache:data | python -m json.tool
+
+# 手动清空开发缓存
+redis-cli -n 1 del faq_search_dev:cache:data faq_search_dev:lock:leader
+
+# 查看当前 leader 是哪个 worker
+redis-cli -n 1 get faq_search_dev:lock:leader
+```

+ 33 - 2
app/config.py

@@ -1,17 +1,38 @@
-"""应用配置:从 .env / 环境变量加载。"""
+"""应用配置:从 .env / 环境变量加载。
+
+通过 APP_ENV 切换不同环境的配置文件:
+    APP_ENV=dev   → 加载 .env.dev(开发环境)
+    APP_ENV=prod  → 加载 .env(生产环境,默认)
+
+启动示例:
+    APP_ENV=dev uvicorn app.main:app --reload                # 开发
+    APP_ENV=prod uvicorn app.main:app --workers 3            # 生产
+    uvicorn app.main:app --workers 3                         # 默认 prod
+"""
+import os
+
 from pydantic_settings import BaseSettings, SettingsConfigDict
 
 
+def _resolve_env_file() -> str:
+    """根据 APP_ENV 选择对应的 env 文件。"""
+    env = os.getenv("APP_ENV", "prod").lower()
+    return ".env.dev" if env == "dev" else ".env"
+
+
 class Settings(BaseSettings):
     """所有可配置项;不允许在源码中硬编码凭据。"""
 
     model_config = SettingsConfigDict(
-        env_file=".env",
+        env_file=_resolve_env_file(),
         env_file_encoding="utf-8",
         case_sensitive=False,
         extra="ignore",
     )
 
+    # 当前运行环境("dev" 或 "prod"),从环境变量 APP_ENV 注入
+    app_env: str = "prod"
+
     # Zendesk 凭据
     zendesk_subdomain: str
     zendesk_email: str
@@ -27,6 +48,16 @@ class Settings(BaseSettings):
     # HTTP
     http_timeout: int = 30
 
+    # Redis(多 worker 共享缓存与 leader 选举)
+    redis_url: str = "redis://127.0.0.1:6379/0"
+    redis_password: str | None = None
+    redis_username: str | None = None
+    redis_key_prefix: str = "faq_search"
+    # leader 锁 TTL(秒);心跳每 leader_lock_ttl/3 续期一次
+    leader_lock_ttl: int = 60
+    # follower 等待 leader 首次写入缓存的最大秒数
+    cache_ready_wait: int = 180
+
     @property
     def zendesk_base_url(self) -> str:
         return f"https://{self.zendesk_subdomain}.zendesk.com/api/v2/help_center"

+ 58 - 18
app/main.py

@@ -12,7 +12,20 @@ from app.exception_handlers import register_exception_handlers
 from app.response import ApiResponse
 from app.routers import search
 from app.schemas import CacheRefreshResult, CacheStatus
-from app.services.cache import faq_cache, scheduler_loop
+from app.services.cache import (
+    WORKER_ID,
+    get_snapshot,
+    leader_heartbeat,
+    refresh_cache,
+    scheduler_loop,
+    try_become_leader,
+    wait_cache_ready,
+)
+from app.services.redis_client import (
+    KEY_LOCK_LEADER,
+    close_redis,
+    release_lock,
+)
 
 logging.basicConfig(
     level=logging.INFO,
@@ -23,27 +36,51 @@ logger = logging.getLogger(__name__)
 
 @asynccontextmanager
 async def lifespan(app: FastAPI) -> AsyncIterator[None]:
-    """启动时初始化缓存 + 启动定时器;关闭时优雅取消。"""
-    logger.info("启动:首次加载 FAQ 缓存…")
-    await faq_cache.refresh()
-
-    task = asyncio.create_task(scheduler_loop(), name="faq-scheduler")
-    logger.info("启动:定时刷新任务已运行")
+    """启动:选举 leader → 仅 leader 刷新缓存与跑定时任务;关闭:释放资源。"""
+    is_leader = await try_become_leader()
+    tasks: list[asyncio.Task[None]] = []
+
+    if is_leader:
+        logger.info("[%s] 当选 leader,开始首次刷新", WORKER_ID)
+        await refresh_cache()
+        tasks.append(
+            asyncio.create_task(scheduler_loop(), name="faq-scheduler")
+        )
+        tasks.append(
+            asyncio.create_task(leader_heartbeat(), name="leader-heartbeat")
+        )
+        logger.info("[%s] 定时任务与 leader 心跳已启动", WORKER_ID)
+    else:
+        logger.info("[%s] 不是 leader,等待缓存就绪…", WORKER_ID)
+        ready = await wait_cache_ready()
+        if ready:
+            logger.info("[%s] 共享缓存已就绪", WORKER_ID)
+        else:
+            logger.warning(
+                "[%s] 等待缓存就绪超时,搜索将以空 sec_ids 降级运行",
+                WORKER_ID,
+            )
 
     try:
         yield
     finally:
-        logger.info("关闭:取消定时任务…")
-        task.cancel()
-        try:
-            await task
-        except asyncio.CancelledError:
-            pass
+        logger.info("[%s] 关闭:取消后台任务…", WORKER_ID)
+        for task in tasks:
+            task.cancel()
+        for task in tasks:
+            try:
+                await task
+            except asyncio.CancelledError:
+                pass
+        if is_leader:
+            # 主动释放 leader 锁,让其他 worker 可以更快接管
+            await release_lock(KEY_LOCK_LEADER, WORKER_ID)
+        await close_redis()
 
 
 app = FastAPI(
     title="Zendesk FAQ Keyword Search",
-    description="封装 Zendesk Help Center 搜索;FAQ sec_ids 每日 0 点刷新缓存。",
+    description="封装 Zendesk Help Center 搜索;FAQ sec_ids 共享缓存(Redis),每日定时刷新。",
     version="1.0.0",
     lifespan=lifespan,
 )
@@ -72,7 +109,7 @@ async def health() -> ApiResponse[dict[str, str]]:
 )
 async def cache_status() -> ApiResponse[CacheStatus]:
     """缓存状态查看。"""
-    snap = faq_cache.snapshot()
+    snap = await get_snapshot()
     return ApiResponse.ok(
         data=CacheStatus(
             sec_ids_count=len(snap["sec_ids"]),
@@ -88,9 +125,12 @@ async def cache_status() -> ApiResponse[CacheStatus]:
     tags=["meta"],
 )
 async def refresh_cache_now() -> ApiResponse[CacheRefreshResult]:
-    """手动触发刷新(运维 / 测试用途)。"""
-    await faq_cache.refresh()
-    snap = faq_cache.snapshot()
+    """手动触发刷新(运维 / 测试用途)。
+
+    使用 force=True 跳过去重锁,确保管理员请求立即生效。
+    """
+    await refresh_cache(force=True)
+    snap = await get_snapshot()
     return ApiResponse.ok(
         data=CacheRefreshResult(
             sec_ids_count=len(snap["sec_ids"]),

+ 3 - 3
app/routers/search.py

@@ -7,7 +7,7 @@ from fastapi import APIRouter, Query
 
 from app.response import ApiResponse, BusinessError, ResponseCode
 from app.schemas import SearchData, SearchRequest
-from app.services.cache import faq_cache
+from app.services.cache import get_snapshot
 from app.services.zendesk_client import ZendeskError, search_articles
 
 logger = logging.getLogger(__name__)
@@ -16,8 +16,8 @@ router = APIRouter(tags=["search"])
 
 
 async def _do_search(req: SearchRequest) -> ApiResponse[SearchData]:
-    """共享的搜索逻辑:使用内存中的 FAQ sec_ids 调用 Zendesk。"""
-    snapshot = faq_cache.snapshot()
+    """共享的搜索逻辑:从 Redis 共享缓存读 sec_ids,调用 Zendesk。"""
+    snapshot = await get_snapshot()
     sec_ids: list[int] = snapshot["sec_ids"]
 
     if not sec_ids:

+ 166 - 64
app/services/cache.py

@@ -1,74 +1,45 @@
-"""FAQ sec_ids 内存缓存 + 每日 0 点定时刷新。"""
+"""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
-from dataclasses import dataclass, field
+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]}"
 
-@dataclass
-class FaqCache:
-    """内存中的 FAQ section 缓存。
-
-    使用不可变赋值(每次刷新整体替换 sec_ids 引用),并发读取安全。
-    写入前用 lock 保护避免并发刷新。
-    """
-
-    sec_ids: list[int] = field(default_factory=list)
-    faq_sections: list[dict[str, Any]] = field(default_factory=list)
-    last_updated_at: datetime | None = None
-    next_refresh_at: datetime | None = None
-    _lock: asyncio.Lock = field(default_factory=asyncio.Lock)
-
-    async def refresh(self) -> None:
-        """从 Zendesk 拉取最新 FAQ sec_ids 替换缓存。"""
-        async with self._lock:
-            logger.info("开始刷新 FAQ 缓存…")
-            try:
-                sec_ids, faq_sections = await list_all_faq_sections()
-            except Exception as exc:
-                # 缓存刷新失败不应让定时任务终止,记录并保留旧值
-                logger.exception("FAQ 缓存刷新失败:%s", exc)
-                return
-
-            # 不可变替换:直接覆盖整体引用,避免读半更新
-            self.sec_ids = list(sec_ids)
-            self.faq_sections = list(faq_sections)
-            self.last_updated_at = datetime.now()
-            self.next_refresh_at = _next_refresh_time(self.last_updated_at)
-            logger.info(
-                "FAQ 缓存刷新完成:sec_ids=%s 条,下一次刷新 %s",
-                len(self.sec_ids),
-                self.next_refresh_at,
-            )
-
-    def snapshot(self) -> dict[str, Any]:
-        """读快照(不持锁)。"""
-        return {
-            "sec_ids": self.sec_ids,
-            "faq_sections": self.faq_sections,
-            "last_updated_at": (
-                self.last_updated_at.isoformat()
-                if self.last_updated_at
-                else None
-            ),
-            "next_refresh_at": (
-                self.next_refresh_at.isoformat()
-                if self.next_refresh_at
-                else None
-            ),
-        }
-
-
-# 全局单例(FastAPI 进程级)
-faq_cache = FaqCache()
+# 空快照(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:
@@ -83,18 +54,128 @@ def _next_refresh_time(now: datetime) -> datetime:
     return today_target + timedelta(days=1)
 
 
-async def scheduler_loop() -> None:
-    """每天 0 点(可配置)刷新一次 FAQ 缓存。
+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 锁。
 
-    采用 sleep(到下一次目标时间) 的简单循环;
-    被 cancel 时优雅退出,无需第三方调度库。
+    锁过期则尝试重新抢锁;抢到继续当 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(
-            "下一次 FAQ 缓存刷新 %s(%.0f 秒后)",
+            "[%s] 下一次 FAQ 缓存刷新 %s(%.0f 秒后)",
+            WORKER_ID,
             next_run.isoformat(),
             wait_seconds,
         )
@@ -103,4 +184,25 @@ async def scheduler_loop() -> None:
         except asyncio.CancelledError:
             logger.info("scheduler_loop 收到取消信号,退出")
             raise
-        await faq_cache.refresh()
+        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

+ 93 - 0
app/services/redis_client.py

@@ -0,0 +1,93 @@
+"""Redis 客户端 + Redis Key 命名空间。
+
+提供进程级单例的 async Redis 连接池,以及统一的 key 前缀生成。
+"""
+from __future__ import annotations
+
+import logging
+from typing import Optional
+
+from redis.asyncio import Redis, from_url
+
+from app.config import settings
+
+logger = logging.getLogger(__name__)
+
+_redis: Optional[Redis] = None
+
+
+def _k(suffix: str) -> str:
+    """构造带前缀的 Redis key。"""
+    return f"{settings.redis_key_prefix}:{suffix}"
+
+
+# 统一管理的 Redis Key
+KEY_CACHE_DATA = _k("cache:data")        # JSON 缓存数据
+KEY_LOCK_LEADER = _k("lock:leader")      # leader 选举锁(值=持有者标识)
+KEY_LOCK_REFRESH = _k("lock:refresh")    # 防并发刷新锁
+
+
+async def get_redis() -> Redis:
+    """获取进程级 Redis 单例(懒初始化)。"""
+    global _redis
+    if _redis is None:
+        _redis = from_url(
+            settings.redis_url,
+            username=settings.redis_username,
+            password=settings.redis_password,
+            encoding="utf-8",
+            decode_responses=True,
+        )
+        logger.info("Redis 连接已建立: %s", settings.redis_url)
+    return _redis
+
+
+async def close_redis() -> None:
+    """关闭 Redis 连接(应用关闭时调用)。"""
+    global _redis
+    if _redis is not None:
+        try:
+            await _redis.aclose()
+        except Exception as exc:  # noqa: BLE001
+            logger.warning("关闭 Redis 连接异常: %s", exc)
+        _redis = None
+
+
+async def try_acquire_lock(
+    key: str, value: str, ttl_seconds: int
+) -> bool:
+    """尝试获取分布式锁(SET NX EX)。成功返回 True。"""
+    r = await get_redis()
+    ok = await r.set(key, value, nx=True, ex=ttl_seconds)
+    return bool(ok)
+
+
+async def renew_lock(key: str, value: str, ttl_seconds: int) -> bool:
+    """续期锁(仅当 value 与持有者一致时才续期,避免误续他人锁)。
+
+    使用 Lua 脚本保证原子性。
+    """
+    r = await get_redis()
+    script = """
+    if redis.call('get', KEYS[1]) == ARGV[1] then
+        return redis.call('expire', KEYS[1], ARGV[2])
+    else
+        return 0
+    end
+    """
+    result = await r.eval(script, 1, key, value, ttl_seconds)
+    return bool(result)
+
+
+async def release_lock(key: str, value: str) -> bool:
+    """释放锁(仅当持有者匹配时才删除)。"""
+    r = await get_redis()
+    script = """
+    if redis.call('get', KEYS[1]) == ARGV[1] then
+        return redis.call('del', KEYS[1])
+    else
+        return 0
+    end
+    """
+    result = await r.eval(script, 1, key, value)
+    return bool(result)

+ 35 - 9
app/services/zendesk_client.py

@@ -6,6 +6,7 @@
 """
 from __future__ import annotations
 
+import asyncio
 import logging
 from typing import Any
 
@@ -16,6 +17,12 @@ from app.config import settings
 logger = logging.getLogger(__name__)
 
 
+# Zendesk 接口请求最大重试次数(包含首次请求)
+MAX_RETRIES = 3
+# 重试退避基数(秒):第 n 次重试等待 RETRY_BACKOFF_BASE * n 秒
+RETRY_BACKOFF_BASE = 1.0
+
+
 class ZendeskError(Exception):
     """Zendesk API 调用失败统一异常。"""
 
@@ -64,17 +71,36 @@ async def list_all_faq_sections(
 
             while url:
                 params: dict[str, Any] = {"per_page": 100, "page": page} if "page" in url else {"per_page": 100}
-                try:
-                    resp = await client.get(url, params=params)
-                except httpx.HTTPError as exc:
-                    logger.error("sections 请求异常 locale=%s: %s", locale, exc)
-                    break
 
-                if resp.status_code != 200:
+                # 最多重试 MAX_RETRIES 次(含首次),全部失败才放弃当前 locale
+                resp: httpx.Response | None = None
+                last_error: str | None = None
+                for attempt in range(1, MAX_RETRIES + 1):
+                    try:
+                        resp = await client.get(url, params=params)
+                        if resp.status_code == 200:
+                            break
+                        last_error = f"status={resp.status_code}"
+                        logger.warning(
+                            "sections 请求失败 locale=%s page=%s attempt=%s/%s %s",
+                            locale, page, attempt, MAX_RETRIES, last_error,
+                        )
+                    except httpx.HTTPError as exc:
+                        last_error = str(exc)
+                        logger.warning(
+                            "sections 请求异常 locale=%s page=%s attempt=%s/%s: %s",
+                            locale, page, attempt, MAX_RETRIES, exc,
+                        )
+                        resp = None
+
+                    # 不是最后一次,等待后再试(指数退避)
+                    if attempt < MAX_RETRIES:
+                        await asyncio.sleep(RETRY_BACKOFF_BASE * attempt)
+
+                if resp is None or resp.status_code != 200:
                     logger.error(
-                        "sections 请求失败 locale=%s status=%s",
-                        locale,
-                        resp.status_code,
+                        "sections 重试 %s 次仍失败 locale=%s page=%s: %s",
+                        MAX_RETRIES, locale, page, last_error,
                     )
                     break
 

+ 1 - 0
requirements.txt

@@ -4,3 +4,4 @@ httpx>=0.27.0
 pydantic>=2.6.0
 pydantic-settings>=2.2.0
 python-dotenv>=1.0.0
+redis>=5.0.0