性能优化策略
概述
覆盖连接复用、并发控制、缓存分层、主机负载分散、资源管理五个维度。本文所有 代码片段均取自实现(文件与函数名可直接检索)。
连接池复用
market_data_service/infra.py::get_http_client() 维护单例 httpx.AsyncClient,
避免每次请求创建/销毁 TCP 连接:
async def get_http_client() -> httpx.AsyncClient:
"""共享 httpx.AsyncClient(延迟初始化 + 事件循环亲和)。"""
global _http_client, _http_client_loop
current_loop = asyncio.get_running_loop()
if (_http_client is None or _http_client.is_closed
or _http_client_loop is not current_loop):
async with _http_client_lock:
...
_http_client = httpx.AsyncClient(
timeout=httpx.Timeout(s.HTTP_TIMEOUT, connect=s.HTTP_CONNECT_TIMEOUT),
limits=httpx.Limits(
max_connections=s.HTTP_MAX_CONNECTIONS,
max_keepalive_connections=s.HTTP_MAX_KEEPALIVE,
keepalive_expiry=s.HTTP_KEEPALIVE_EXPIRY,
),
headers=EM_HEADERS,
follow_redirects=True,
)
_http_client_loop = current_loop
return _http_client
双检锁 + 事件循环亲和(_http_client_loop):绑定在旧事件循环上的客户端在新
循环里使用会抛 RuntimeError(热重载/测试多循环场景),检测到 loop 变更即重建。
优雅关闭走 close_http_client()。
| 参数 | 默认值 | 配置项(core/config.py) |
|---|---|---|
| 请求总超时 | 10s | HTTP_TIMEOUT |
| 连接超时 | 5s | HTTP_CONNECT_TIMEOUT |
| 最大连接数 | 20 | HTTP_MAX_CONNECTIONS |
| 保活连接数 | 10 | HTTP_MAX_KEEPALIVE |
| 保活超期 | 30s | HTTP_KEEPALIVE_EXPIRY |
并发获取
互不依赖的数据并发获取,示例为 _fetch_from_tdx() 的 sentiment 分支
(线程池执行同步 TDX/DB 调用,再 gather 汇合,外层带 _SOURCE_GATHER_BUDGET 超时):
indices_task = loop.run_in_executor(executor, tdx_indices)
breadth_task = loop.run_in_executor(executor, tdx_market_breadth)
indices, breadth = await asyncio.wait_for(
asyncio.gather(indices_task, breadth_task),
timeout=_SOURCE_GATHER_BUDGET,
)
if not breadth:
breadth = await loop.run_in_executor(executor, db_market_breadth)
线程池配置
# 8 线程:外部源偶发挂起时避免占满线程池,阻塞纯 DB 查询(如 hot_stocks)
executor = ThreadPoolExecutor(max_workers=8, thread_name_prefix="mkt")
- 同步 TDX(TCP)/ DB 调用放入线程池,避免阻塞 asyncio 事件循环
- 8 线程是为了故障隔离而非吞吐:外部源挂起会占满线程,需给
hot_stocks一类的纯 DB 查询留出并行余量
编号主机负载分散
东财 17.push2 / 40.push2 等编号主机统一定义在 em_paginated.py(单一来源原则),
em_client.py 从此处导入,不再重复定义:
# em_paginated.py — push2 主机池唯一定义源
PUSH2_HOSTS = [
"push2.eastmoney.com",
"17.push2.eastmoney.com",
"29.push2.eastmoney.com",
"40.push2.eastmoney.com",
"79.push2.eastmoney.com",
"91.push2.eastmoney.com",
]
PUSH2_HIS_HOSTS = [
"push2his.eastmoney.com",
"7.push2his.eastmoney.com",
"33.push2his.eastmoney.com",
"63.push2his.eastmoney.com",
"91.push2his.eastmoney.com",
]
em_paginated.py::_next_push2_host() 对 6 个 push2 节点做轮转(_push2_host_idx
自增取模),避免单节点过载触发风控。
em_client.py 在导入时将裸域名 push2.eastmoney.com 从镜像池中剔除
(_EASTMONEY_PUSH2_HOSTS = [h for h in PUSH2_HOSTS if h != "push2.eastmoney.com"])——
裸域名是首选主机,编号节点作为失败后的镜像池,由 HostFallbackManager 管理:
- 主机失败进入指数退避冷却:30s → 60s → … → 上限 300s,成功即复位(
COOLDOWN_MS/COOLDOWN_MAX_MS) - 冷却期内的主机被跳过,请求自动切换到池中其它节点
push2delay(延时行情域)纳入同一套冷却管理,避免"单主机池不判定"导致每次失败都白等完整重试链
东财涨跌分布分阶段请求
providers_eastmoney.py::em_market_breadth() 的分段依赖 push2 filter 参数,而该参数
已被东财服务端停用(2026-09-22 实测),故请求分三阶段、按需推进:
# 阶段一:2 路探针(全量 + 上涨),验证 filter 是否生效
data_total, data_up = await asyncio.gather(
em_fetch_with_retry(url, base_params),
em_fetch_with_retry(url, {**base_params, "filter": "(f3>0)"}),
)
# filter 失效时 up == total(全量总数)→ 立即放弃,不再发后续 15 路
if up >= total:
logger.warning("东财区间统计 filter 参数已失效…交降级链")
return None
# 阶段二:跌/涨停/跌停三路并发
# 阶段三:13 段区间并发,用 return_exceptions=True 做局部失败降级
# (任一段失败则整体放弃 bins,仅返回家数,不阻断主功能)
阶段一的意义:filter 失效时把 17 路请求降到 2 路。东财有风控,无效请求同样计入
HostFallbackManager 主机冷却,会牵连龙虎榜/两融等东财独有数据。最终结果仍要过
distribution_counts_consistent(家数/分段/总数自洽)兜底校验,filter 部分失效
同样判不可用。
请求去重
同一 data_type 在同一交易日(休市期间)仅调用一次外部 API,实现位于 infra.py:
_request_log: dict[str, str] = {} # {data_type: "2026-09-21"}
def already_fetched_today(data_type: str) -> bool:
return _request_log.get(data_type) == today_str()
def mark_fetched(data_type: str): ... # 仅在"有效数据成功取得"后调用
- 去重仅在休市时生效(交易时段数据持续变化,必须实时拉取)
- 实时编排超时的路径不写
mark_fetched、不落快照,语义等同"获取失败",下次请求会重试
静态数据 lru_cache
不变的参考数据用 @lru_cache(maxsize=1) 缓存,避免每次请求重复构建
(market_data_service/infra.py):
@lru_cache(maxsize=1)
def index_code_map() -> dict[str, str]: ... # 指数代码 → 名称
@lru_cache(maxsize=1)
def index_tencent_codes() -> dict[str, str]: ... # 指数代码 → 腾讯带市场前缀代码
@lru_cache(maxsize=1)
def index_secids() -> str: ... # 东财 ulist 用 secids 串
@lru_cache(maxsize=1)
def breadth_sample_codes() -> tuple[str, ...]: ... # 涨跌分布抽样代表股(50 只)
@lru_cache(maxsize=1)
def em_all_market_fs() -> str: ... # 东财 clist 全市场筛选条件
东财响应缓存另在 em_client.py 用 @lru_cache(maxsize=_EM_CACHE_SIZE)(64 项)。
DB 会话与可用性探测
会话统一由 core/database.py 提供,两条通道:
def get_db() -> Generator[Session | None, None, None]: ... # FastAPI 依赖注入
@contextmanager
def get_db_session() -> Generator[Session | None, None, None]:
"""获取数据库会话(可在 with 语句中使用)。"""
if SessionLocal is None:
yield None
return
db = SessionLocal()
try:
yield db
finally:
db.close()
db_available() 做真实连接探测(仅检查 engine is None 是假阳性:MySQL 宕机时
engine 依然存在,接口会以 500 崩溃而非优雅降级),探测结果缓存 5s(_DB_PROBE_TTL)。
未配置数据库时 SessionLocal 为 None,get_db_session() 产出 None,调用方按
mock 模式处理,不抛异常。
性能指标(设计目标与实测)
| 指标 | 值 | 实现方式 |
|---|---|---|
| 内存缓存命中 | < 5ms | 进程内 mkt_{data_type} 缓存,交易时段 TTL 10~30s |
| 单次请求外部调用上限 | 6 次 | RETRY_BUDGET_MAX(防止雪崩式重试) |
| 单源并发汇合预算 | 10s | _SOURCE_GATHER_BUDGET(TDX 指数 + 涨跌统计并发) |
| 实时源编排总预算 | 15s 硬上限 | _LIVE_FETCH_BUDGET,超时即降级 DB 历史兜底 |
| TDX 实时取数 | 无 HTTP 往返 | TCP 长连接(自实现协议),不封 IP |
| 腾讯 K 线冷启动 | 实测 avg ~162ms | 共享连接池 + 前复权单请求 |
| 全源失败 | 立即返回 | fox 持久层历史兜底 / 空结构,不做额外等待 |
| 内存占用 | 稳定 | TTL 自动过期 + lru_cache 限量 |