mirror of
https://ghfast.top/https://github.com/aeroxw/easy_tdx_max.git
synced 2026-09-12 22:44:22 +08:00
fix(web): 单飞改用轮询解耦,预热超时不再波及端点请求
问题:上一版单飞用 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 兜底
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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")
|
||||
|
||||
Reference in New Issue
Block a user