8.4 KiB
市场数据同步性能与完整性检查设计
1. 目标与边界
本设计保留现有“逐股获取六年 qfq 完整快照”的业务口径,通过 8 路受控并发和 PostgreSQL 连接/写入批量化缩短日常同步。另提供不访问 Tushare、不修改市场事实的手动完整性检查。
不引入按交易日增量下载、自动完整检查、自动修复或通用任务队列。选股仍只消费已完成且 覆盖率达标的市场同步批次。
2. 任务拆分与依赖
父任务只负责需求、跨子任务契约和最终集成,不直接实现产品代码。
08-11-market-sync-concurrency-storage- 交付 8 worker、共享频控、连接池、批量 item 写入和集合覆盖率。
- 无子任务依赖,必须最先完成。
08-11-market-integrity-check-api- 依赖第 1 项提供的池化 PostgreSQL 基础设施和 advisory lock 语义。
- 交付检查领域模型、迁移、只读比较、HTTP 触发和轮询。
08-11-market-integrity-check-web- 依赖第 2 项冻结 HTTP 字段和状态值。
- 交付
/sync页面、导航、触发、轮询和报告展示。
3. 同步执行模块
3.1 外部接口
SyncMarketData.execute(command) -> SyncBatchSummary 保持不变。并发数、请求协调器和仓储
由组合根注入,CLI 调用方不需要了解线程、频控或连接池细节。
3.2 单股 worker
把当前会修改共享 failures / totals 的 _process_bar 改为返回不可变
SyncItemOutcome:股票代码、状态、WriteResult、fingerprint 和可选安全错误。
每个 worker 内仍严格执行:
- 获取该股票六年 qfq;
- 读取旧正式 CSV 并比较重叠指纹;
- 写临时 CSV;
- 在独立 PostgreSQL 事务中幂等 upsert;
- 数据库提交后原子发布 CSV;
- 返回 outcome。
主线程通过 as_completed 聚合 outcome、更新进度和每 100 项批量写入同步审计。单股异常只
生成失败 outcome;未发布的临时文件必须清理。结果顺序不影响计数或最终状态。
3.3 请求协调器
请求协调器是 Tushare 适配器内部的深模块,接口只暴露“执行一个真实供应商调用”。实现隐藏:
- 正常请求保持最多 8 路并发,不用一个全局固定间隔把真实调用重新串行化;
- 保留现有可配置的单股请求间隔作为温和节流;
- “访问频繁/请稍后/超过频率/too many requests/429/403”等频控分类;
- 频控时共享
cooldown_until,默认按 60、120、180 秒增长并设置上限; - 非频控临时错误沿用有界重试和普通退避;
- 注入 monotonic clock / wait 实现确定性单元测试;
- 日志只包含方法名、尝试次数和等待时间,不包含 token 或供应商原始响应。
为保留 Tushare qfq 计算,继续调用 ts.pro_bar(api=coordinated_client, adj="qfq", ...)。
coordinated_client.daily 与 coordinated_client.adj_factor 在原始异常仍可见的位置经过请求
协调器;pro_bar 内部重试降为一次,避免隐藏重试绕开全局频控。
3.4 PostgreSQL
- 增加
psycopg_pool.ConnectionPool依赖,池最大连接数与 worker 数匹配并留出控制连接。 - CLI 用上下文管理器打开/关闭池;HTTP 进程使用按数据库 URL 缓存的池并在进程退出时关闭。
- worker 的 bar upsert 从池中借一个连接并保持单股票事务,不在线程间共享 connection。
- 主线程调用
record_items(batch_id, outcomes)分批 upsertmarket_sync_item。 - 仓储新增
count_valid_stocks(target_trade_date),用一条集合 SQL 计算 active 股票中同时存在 bar 和 daily_basic 的数量;删除同步编排对逐股has_*的依赖。 - advisory lock 的连接在整个批次期间保持借出,不和 worker 连接混用。
4. 完整性检查模块
4.1 语义
“完整性检查”验证 PostgreSQL 事实与已发布 CSV 快照是否一致,并验证 CSV 可以按现有领域规则 解析。它不判断 Tushare 是否已发布某日数据,也不能发现 PostgreSQL 和 CSV 同时缺少但供应商 实际存在的数据。
检查范围:
- 当前 active 股票主数据与
stock-basic/current.csv; - 最近成功同步窗口内的每股 qfq bar 与
bars/<ts_code>.csv; - 同一窗口内 PostgreSQL daily_basic 与
daily-basic/<YYYY>/<YYYYMMDD>.csv; - 缺失、多余、无法解析、重复键、字段内容不一致和窗口越界。
4.2 独立持久化模型
新增表而不复用 market_sync_batch:
market_integrity_checkid,status,window_start,window_endtarget_count,checked_count,issue_counterror_type,error_message,created_at,finished_at
market_integrity_issuecheck_id,item_kind,item_key,issue_type,message,created_at- 复合主键包含可稳定区分同一对象多个问题的序号或 issue key。
状态固定为:running、passed、issues_found、failed。issues_found 表示检查完整执行但数据
存在问题,不等于检查任务执行失败。
4.3 只读比较
RunMarketIntegrityCheck 是外部应用接口:
prepare() -> IntegrityCheckRun原子创建 running 记录;已有 running 时返回冲突。execute(check_id) -> None获取与同步相同的 advisory lock,执行比较并收敛终态。get(check_id, page, page_size)和get_latest(...)提供轮询读模型。
PostgreSQL adapter 使用有序 server-side cursor,按股票或交易日流式产生一组记录;CSV adapter 按同样 key 顺序读取。应用模块做 merge comparison,一次只保留一个股票或一个交易日的数据, 不把约 700 万行全量装入内存。
检查持有与同步相同的 advisory lock,保证 DB 与正式 CSV 在比较期间不会被同步修改。它只向 检查表写进度和问题;市场事实表和正式 CSV 不发生写入。
4.4 HTTP
路由归属 market_data presentation,并由 /api/v1/market-data 挂载:
POST /integrity-checks→202 Accepted,返回 check id、status、窗口;GET /integrity-checks/latest→ 最近一次检查或no_data;GET /integrity-checks/{check_id}?page=&page_size=→ 进度、终态和分页问题。
HTTP 使用仓库已有的 FastAPI BackgroundTasks 模式。后台任务出现未捕获异常时必须把检查记录
收敛为 failed。进程重启可能中断 in-process task;下一次触发允许显式把失去 advisory lock
且超过超时阈值的 running 记录标记为 failed,避免永久阻塞。该恢复只修改检查元数据。
5. Web 模块
启用现有“同步任务”导航并新增 /sync 路由,使用独立 features/sync 垂直切片:
- API 类型与函数;
- 最新检查 query、单次检查 polling query、触发 mutation;
- 页面展示检查说明、报告只读提示、最近窗口、进度、状态和分页问题;
- running 时禁用重复触发并按固定间隔轮询;终态停止轮询并刷新 latest;
409显示已有检查正在运行,503显示存储不可用,其他错误保留可重试入口。
页面不把服务器状态复制到 Zustand。问题明细最少显示对象类型、key、问题类型和安全说明。
6. 兼容、部署与回滚
- 新迁移为
0003_market_integrity_checks,只新增检查表/索引,不修改事实表。 - 配置默认 worker 从当前未生效的 4 调整为 8;新增频控退避配置时同步
.env.example、 Compose 和市场数据运维文档。 - CLI 参数和退出码保持兼容;旧批次和旧 CSV 无需迁移。
- 回滚应用版本时新增检查表可暂留;数据库 downgrade 只删除检查表,不触碰市场事实。
- 若并发上线后供应商频控持续恶化,可通过环境变量把 worker 降为 1,无需回滚代码。
7. 验证策略
- 确定性并发测试证明 active worker 不超过配置值,结果汇总不依赖完成顺序。
- fake clock/condition 测试证明频控触发共享冷却,普通错误不冻结其他 worker。
- PostgreSQL 集成测试验证池化并发 upsert、批量 item、集合覆盖率和事务回滚。
- 完整性检查 fixture 覆盖 passed、缺 CSV、缺 DB、内容不一致、非法 CSV、检查冲突和批次异常。
- HTTP 测试锁定 202、409、分页与状态契约;Web 测试锁定触发、轮询停止、报告只读提示和错误态。
- 无网络 5002 股票调度基准锁定线程与仓储调用数量;生产约 30 分钟目标以部署日志验收。