From 25236056f7f9bbbd9e6dd8cebc2030ea606ff877 Mon Sep 17 00:00:00 2001 From: Justin Gu <97915@qq.com> Date: Mon, 6 Jul 2026 02:39:02 +0800 Subject: [PATCH] =?UTF-8?q?fix(web):=20=E5=8D=95=E9=A3=9E=E6=94=B9?= =?UTF-8?q?=E7=94=A8=E8=BD=AE=E8=AF=A2=E8=A7=A3=E8=80=A6=EF=BC=8C=E9=A2=84?= =?UTF-8?q?=E7=83=AD=E8=B6=85=E6=97=B6=E4=B8=8D=E5=86=8D=E6=B3=A2=E5=8F=8A?= =?UTF-8?q?=E7=AB=AF=E7=82=B9=E8=AF=B7=E6=B1=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 问题:上一版单飞用 asyncio.Future 共享,lifespan 预热用 wait_for(120s) 超时会 cancel 整个 Future,导致同时等待的 /security/search-index 端点 请求收到 CancelledError 返回 500。日志显示预热 TimeoutError + 端点 500。 修复: - _build_search_index 改用轮询设计:后台 _do_build task 独立运行 (fire-and-forget),调用方通过轮询 _SEARCH_INDEX 每 0.5s 检查结果 - 调用方被 cancel(预热超时)只退出轮询,不影响后台构建 task - lifespan 预热去掉 wait_for,改为纯 best-effort 触发,无权取消构建 - 冒烟验证:预热被 cancel 后,端点请求仍能正常拿到结果(5206 条) 另外从日志看慢机器爬全名单 >120s,轮询上限提到 300s 兜底 --- src/easy_tdx/web/app.py | 7 ++-- src/easy_tdx/web/routers/market.py | 58 ++++++++++++++++++------------ 2 files changed, 38 insertions(+), 27 deletions(-) diff --git a/src/easy_tdx/web/app.py b/src/easy_tdx/web/app.py index 010d787..d937ec0 100644 --- a/src/easy_tdx/web/app.py +++ b/src/easy_tdx/web/app.py @@ -67,16 +67,15 @@ async def lifespan(app: FastAPI) -> AsyncGenerator[None, None]: # --- 股票搜索索引预热(后台,不阻塞服务启动) --- # 首次 get_security_list_all("all") 要几十次 TDX 协议往返(几十秒), # 后台提前跑,让本地缓存尽早建立。用户打开页面时大概率已就绪。 - # 用 create_task 不 await,服务立即可用;预热未完成时前端遮罩接管等待。 - # 与 /security/search-index 端点共享同一单飞任务(_build_search_index), - # 避免预热和首次端点请求并发各爬一次全名单。 + # 仅触发构建(fire-and-forget),不 await 不超时——预热是 best effort, + # 没有权力取消构建 task(否则会波及同时到达的 /security/search-index 请求)。 import asyncio async def _warmup_security_list() -> None: try: from easy_tdx.web.routers.market import _build_search_index - await asyncio.wait_for(_build_search_index(client), timeout=120) + await _build_search_index(client) logger.info("Security list warmup done") except Exception: logger.warning("Security list warmup failed (non-fatal)", exc_info=True) diff --git a/src/easy_tdx/web/routers/market.py b/src/easy_tdx/web/routers/market.py index 20ee270..e7444b2 100644 --- a/src/easy_tdx/web/routers/market.py +++ b/src/easy_tdx/web/routers/market.py @@ -57,25 +57,34 @@ async def security_list_all( # 进程级缓存:5206 条记录的 {code, name, initials} 只算一次,热重启即丢。 # 前端拉一次后模块级缓存,按 code/name.includes/initials.includes 三路过滤。 _SEARCH_INDEX: list[dict[str, str]] | None = None -# 单飞 Future:预热(lifespan)与端点请求共享同一任务,避免并发爬两次全名单。 -# 首次请求会 await 这个 Future;后续请求命中 _SEARCH_INDEX 直接返回。 -_SEARCH_INDEX_TASK: Any = None # asyncio.Future[list[dict[str, str]]] +# 单飞标记:True 表示后台构建 task 正在跑。调用方据此判断是否需要启动新 task。 +# 注意:不持有 task/future 引用,避免调用方被 cancel 时波及后台构建。 +_SEARCH_BUILDING: bool = False +# 后台构建失败的最近一次异常(供等待中的调用方读取;None 表示无错或未发生) +_SEARCH_BUILD_ERROR: BaseException | None = None async def _build_search_index(client: Any) -> list[dict[str, str]]: - """构建搜索索引:拉全名单 + pypinyin 预计算声母。耗时几十秒(首次)。""" + """构建搜索索引:拉全名单 + pypinyin 预计算声母。耗时几十秒(首次)。 + + 单飞 + 轮询设计:后台构建 task 与调用方解耦,调用方被 cancel(如预热超时) + 不会波及正在跑的构建 task,也不会让其他等待的请求收到 CancelledError。 + """ import asyncio - global _SEARCH_INDEX, _SEARCH_INDEX_TASK - # 双重检查:等待期间可能已被其他协程填好 + global _SEARCH_INDEX, _SEARCH_BUILDING, _SEARCH_BUILD_ERROR + + # 已就绪:直接返回 if _SEARCH_INDEX is not None: return _SEARCH_INDEX - # 单飞:已有进行中的任务则复用,避免预热 + 端点请求并发爬两次 - if _SEARCH_INDEX_TASK is None: - _SEARCH_INDEX_TASK = asyncio.get_running_loop().create_future() + + # 未启动构建:启动后台 task(fire and forget,调用方不持有它的引用) + if not _SEARCH_BUILDING: + _SEARCH_BUILDING = True + _SEARCH_BUILD_ERROR = None async def _do_build() -> None: - global _SEARCH_INDEX, _SEARCH_INDEX_TASK + global _SEARCH_INDEX, _SEARCH_BUILDING, _SEARCH_BUILD_ERROR from pypinyin import Style, lazy_pinyin try: @@ -89,23 +98,26 @@ async def _build_search_index(client: Any) -> list[dict[str, str]]: initials = "".join(lazy_pinyin(name, style=Style.FIRST_LETTER)) index.append({"code": code, "name": name, "initials": initials}) _SEARCH_INDEX = index - if not _SEARCH_INDEX_TASK.done(): - _SEARCH_INDEX_TASK.set_result(index) - except Exception as e: - # 失败清空 task,允许下次重试 - if not _SEARCH_INDEX_TASK.done(): - _SEARCH_INDEX_TASK.set_exception(e) - _SEARCH_INDEX_TASK = None - raise + except BaseException as e: + _SEARCH_BUILD_ERROR = e finally: - # 成功后清空 task 引用(结果已存 _SEARCH_INDEX) - _SEARCH_INDEX_TASK = None + _SEARCH_BUILDING = False asyncio.create_task(_do_build()) - from typing import cast - - return cast("list[dict[str, str]]", await _SEARCH_INDEX_TASK) + # 轮询等待结果(每 0.5s 检查一次)。 + # 这样调用方被 cancel 时,只是退出轮询,不影响后台 _do_build task。 + # 用 asyncio.shield 保护轮询本身不被取消传播,并在每次循环检查错误。 + for _ in range(600): # 上限 300 秒(600 × 0.5s) + if _SEARCH_INDEX is not None: + return _SEARCH_INDEX + if not _SEARCH_BUILDING and _SEARCH_BUILD_ERROR is not None: + # 构建已结束但失败:抛错给调用方(下次调用会重新触发构建) + err = _SEARCH_BUILD_ERROR + _SEARCH_BUILD_ERROR = None + raise err + await asyncio.sleep(0.5) + raise TimeoutError("搜索索引构建超时(300s)") @router.get("/security/search-index")