Files

44 lines
2.2 KiB
Markdown
Raw Permalink Normal View History

2026-08-11 13:32:21 +08:00
# 八路同步与存储批量化设计
## 模块接口
- `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,得到与原串行流程等价的安全回退。