Files
2026-08-12 09:45:27 +08:00

6.6 KiB
Raw Permalink Blame History

技术设计:选股执行性能优化

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 能力。它内部增加三个私有阶段:

  1. load_histories:按股票分块从 reader 读取历史;
  2. evaluate_histories:使用固定数量的线程 worker 调用现有 EvaluateZhixingB1.execute_history;
  3. 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 在一个事务内:

  1. 按分块股票代码删除该 run/股票已有 signal,保证重试幂等;
  2. 使用 psycopg connection cursor 的 executemany upsert 全部 selection_run_item;
  3. 使用同一批量 API 插入全部独立 selection_signal;旧 fake connection 没有 cursor 时保留逐条 execute fallback;
  4. 事务成功后返回。

原 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 采用 用户确认的初始并发度,后续以基准调整