Files
zhixing-system/.trellis/tasks/archive/2026-08/08-11-market-sync-concurrency-storage/design.md
T
2026-08-11 13:32:21 +08:00

2.2 KiB
Raw Blame History

八路同步与存储批量化设计

模块接口

  • SyncMarketData.execute(command):外部接口保持不变,内部使用固定大小 executor。
  • SyncItemOutcome:worker 返回值,封装 result/fingerprint/failure;不暴露线程实现。
  • Tushare 请求协调器:适配器内部接口,执行一个真实 client 方法并隐藏间隔、重试和共享冷却。
  • MarketDataRepository.record_items(...):批量持久化 outcomes。
  • MarketDataRepository.count_valid_stocks(...):集合式覆盖率接口。

并发流

stock master 和 daily_basic 保持主线程顺序执行。bar 阶段把股票提交给最多 8 个 worker;每个 worker 完成获取、CSV staging、数据库事务和 CSV 发布后返回。主线程 as_completed 聚合并每 100 项刷新 审计和日志。批量审计写失败是批次级错误,不回滚已提交事实。

Tushare 频控

由 coordinated client 代理 pro_bar 实际调用的 daily / adj_factor:

  1. 请求前检查共享 cooldown;正常情况下允许最多 8 路并发;
  2. 调用原始 token client;
  3. 频控异常扩大共享 cooldown,唤醒/阻塞所有等待 worker;
  4. 普通临时错误只退避当前调用;
  5. 到达重试预算后抛出安全 TushareSourceError。

现有可配置请求间隔保留为单股成功后的温和节流,不对所有真实调用强加全局串行间隔。 pro_bar(..., retry_count=1) 避免 SDK 在看不到共享 gate 的位置自行重试。测试注入 fake clock、wait 和 client,不使用真实睡眠或网络。

PostgreSQL

依赖调整为 Psycopg pool extra。CLI 在一次执行期间打开池并在退出时关闭;池最大连接数至少覆盖 8 个 worker、一个主线程写连接和 advisory lock 连接。connection 不跨线程共享。

record_items 使用 executemany 或 COPY/staging 一次写一批。count_valid_stocks 从 active 股票 连接/EXISTS 两张目标日事实表后 count,一次返回 valid count。

失败与回退

  • executor 创建或批量审计失败记录 batch 错误并收敛状态。
  • worker 异常必须转换为 outcome,不能让 future 异常跳过进度。
  • worker 数可降为 1,得到与原串行流程等价的安全回退。