6.6 KiB
技术设计:选股执行性能优化
1. 设计目标与边界
本次只优化 selection bounded context 的执行数据流,不改变 ZhixingB1Strategy
的公式实现、StockHistory 的业务语义和现有 HTTP 结果契约。
保留以下外部事实:
- PostgreSQL 是行情和选股最终结果的事实源;
- 选股运行仍由
POST /api/v1/selection/runs创建并返回run_id; - 每个股票仍有一个
SelectionRunItem,每类命中仍有一个独立SelectionSignal; - 单股评估异常继续隔离,不中断整个运行;
(strategy, target_trade_date)的运行 claim、重跑保护和 advisory lock 保持不变。
不在本次设计中加入 Redis、外部任务队列、进程级任务恢复或新的前端进度协议。
2. 模块与接缝
2.1 应用模块
继续由 RunZhixingB1 作为深模块,对 presentation 暴露 prepare/execute/query
能力。它内部增加三个私有阶段:
load_histories:按股票分块从 reader 读取历史;evaluate_histories:使用固定数量的线程 worker 调用现有EvaluateZhixingB1.execute_history;record_items:按结果分块提交 PostgreSQL。
公式策略本身仍只接收一个 StockHistory,避免把数据库、连接池和并发细节带入
domain。
2.2 批量读取接口
为兼容已有单股调用者和测试 fake,在 SelectionUniverseReader 旁增加可选的批量
历史读取扩展协议;生产适配器实现该扩展,执行器运行时优先探测批量方法,缺失时
保留原单股 fallback。批量历史读取能力为:
load_histories(
stocks: Sequence[SelectionStock],
target_trade_date: date,
) -> tuple[StockHistory, ...]
load_history(ts_code, target_trade_date) 保留给单股调用者和兼容测试;批量执行路径
不再调用它。
PostgresMarketDataReader.load_histories 使用一条参数化 SQL 读取一个分块:
bar.ts_code = ANY(%s);bar.source_adj = 'qfq';bar.trade_date <= %s;- 按
bar.ts_code, bar.trade_date升序返回; - 当前 B1 路径只读取股票名和 OHLCV,不再为历史每一行 LEFT JOIN
market_daily_basic;目标日 basic 的完整性仍由执行源查询检查。
读取结果按股票代码分组,缺失代码返回空 StockHistory,由现有策略状态矩阵映射为
missing_target_bar 或 insufficient_history。不会改变目标日截断、qfq 和排序契约。
分块大小由 selection_batch_size 控制,默认 200;每个分块读取完成、评估完成并写入
后才释放历史对象,避免全市场历史同时驻留内存。
2.3 有界 PostgreSQL 资源
新增 selection infrastructure 的轻量资源 owner,内部持有一个
psycopg_pool.ConnectionPool:
max_size = selection_max_workers + 2,默认 6;- reader 和 run repository 共享同一个 pool;
- pool 在进程内按
(database_url, max_workers)缓存; - 首次使用时 open,进程退出时 close;
- 测试通过构造函数注入 fake pool/connection。
这沿用市场数据模块现有的 pool 生命周期模式,不让每个 HTTP 请求或每只股票拥有 独立 pool。读写 adapter 只借用短生命周期连接,业务事务仍由 adapter 控制。
2.4 批量写入接口
为兼容已有单项调用者和测试 fake,在 SelectionRunStore 旁增加可选的批量写入
扩展协议:
record_items(run_id: str, items: Sequence[SelectionRunItem]) -> None
PostgresSelectionRunRepository.record_items 在一个事务内:
- 按分块股票代码删除该 run/股票已有 signal,保证重试幂等;
- 使用 psycopg connection cursor 的
executemanyupsert 全部selection_run_item; - 使用同一批量 API 插入全部独立
selection_signal;旧 fake connection 没有 cursor 时保留逐条 execute fallback; - 事务成功后返回。
原 record_item 保留为单项兼容 wrapper,并委托给 record_items([item]);应用批量
路径不再调用它。新 run 的每个写入分块只提交一次事务,单股公式异常仍在应用层被
转成 data_error 后进入该批次。
2.5 四 worker 执行模型
RunZhixingB1 构造时接收 max_workers=4,也可由 Settings 注入。每个 history 分块
使用一个长期存在的 ThreadPoolExecutor(max_workers=4) 评估;worker 不直接写库。
选择线程而不是立即引入进程池的原因:
- 历史读取已经在应用线程按分块完成,无需跨进程复制数据库连接;
- pandas/numpy 的一部分计算可以释放 GIL;
- 线程共享只读
StockHistory,实现和回滚成本较低; - 若基准证明公式 CPU/GIL 成为主瓶颈,后续可以在同一 evaluator 接缝替换为进程 worker,不影响 reader/store 契约。
每个分块保持如下状态流:
读取分块 → 4 worker 评估 → 聚合计数 → 一次批量写入 → 处理下一分块
评估完成顺序不作为业务契约;查询端仍按股票代码和公式优先级稳定排序。
3. 配置与可观测性
新增 Settings 字段:
selection_max_workers: int = 4,环境变量ZHIXING_SELECTION_MAX_WORKERS;selection_batch_size: int = 200,环境变量ZHIXING_SELECTION_BATCH_SIZE。
执行结束时写一条安全的汇总日志,包含:run id、股票数、历史行数、分块数、worker 数以及 read/evaluate/persist 的 wall-clock 秒数。不得写入连接串、token、行情详情 或完整异常堆栈。
4. 兼容性与回滚
- 不新增或修改数据库表、索引和 HTTP 字段;无需 migration。
- 如果批量读取/写入出现问题,可暂时将
selection_batch_size=1,保留相同接口并 回到单股票批次;max_workers=1可关闭并发以定位问题。 - 通过 golden、应用层多分类测试和结果存储测试保证公式结果不漂移。
- 任何 pool 获取失败都按现有 storage error 语义处理,不把半批次伪装成成功。
5. 关键取舍
| 选择 | 本次决定 | 原因 |
|---|---|---|
| Redis | 不引入 | 当前问题是 PostgreSQL 往返和串行执行,已有持久化事实源足够 |
| 全市场一次读取 | 不采用 | 六年历史乘以全股票池会增加内存峰值 |
| 分块读取 | 采用,默认 200 | 控制内存并保留中间进度/失败隔离 |
| 每股事务 | 不采用 | 事务数量随股票数线性增长 |
| 分块事务 | 采用 | 减少提交次数,同时保留可控的部分进度 |
| 无界并发 | 不采用 | 可能耗尽 PostgreSQL 连接和内存 |
| 4 个线程 worker | 采用 | 用户确认的初始并发度,后续以基准调整 |