Files
2026-08-06 22:58:27 +08:00

135 lines
10 KiB
Markdown
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# Tushare PostgreSQL 同步设计
## 1. 设计目标
建立一个不依赖 FastAPI 生命周期的一次性市场数据同步用例。它负责把 Tushare 当前沪深非 ST 股票池、六年 qfq 日线和按交易日归属的每日指标写入 PostgreSQL,同时发布可审计、可恢复的 CSV 落地快照,并在局部失败时给后续策略提供明确的数据新鲜度和覆盖率。
本任务保持为一个 Trellis 任务:领域模型、数据库迁移、同步编排和 Compose Job 共用同一套数据契约与端到端验收,拆成独立子任务会在 schema 和运行契约尚未落地时制造额外依赖管理。
## 2. 模块边界
新增 `modules/market_data` bounded context,按项目约定分层:
```text
modules/market_data/
├── domain/ # 股票、日线、每日指标、窗口、批次状态、端口
├── application/ # 初始化、日常同步、失败重试和覆盖率编排
├── infrastructure/ # Tushare、CSV、Psycopg 适配器
└── presentation/ # 一次性 CLI 入口
```
- `domain` 定义值对象、规范化/指纹规则和外部端口,不导入 FastAPI、Psycopg、pandas 或 Tushare。
- `application` 只依赖端口,暴露一个深接口 `SyncMarketData.execute(command) -> SyncBatchSummary`;初始化、日常运行和失败重试由 command 参数表达,不复制三套流程。
- `infrastructure` 封装 Tushare SDK、原子文件发布、Psycopg 事务与 SQL。
- `presentation` 解析 CLI 参数、组装依赖并返回适合 cron 判断的退出码;不增加 HTTP 触发接口。
- 配置仍由 `bootstrap.config.Settings` 统一读取,业务代码不直接访问环境变量。
## 3. 领域与数据契约
### 3.1 标识与时间窗口
- 使用 Tushare `ts_code` 作为规范股票代码;数据库逻辑唯一键为 `(ts_code, trade_date)`。
- CLI 可显式传入目标交易日;未传入时通过 Tushare 交易日历解析不晚于当前日期的最近开市日。周末、节假日和重复触发因此指向同一目标日期并保持幂等。
- `window_start` 为目标交易日向前移动六个日历年的同月同日,边界包含在同步窗口中;数据库和 CSV 清理 `< window_start` 的数据。
### 3.2 PostgreSQL schema
迁移创建以下业务表;表名最终可按项目迁移命名规范微调,但字段职责不变:
| 表 | 主键/唯一键 | 主要职责 |
| --- | --- | --- |
| `market_stock` | `ts_code` | 当前股票主数据、上市状态、名称、市场、交易所和刷新时间 |
| `market_daily_bar` | `(ts_code, trade_date)` | qfq OHLC、昨收、涨跌、成交量和成交额 |
| `market_daily_basic` | `(ts_code, trade_date)` | 换手率、量比、估值、股本与市值等按日指标 |
| `market_sync_batch` | `id` | 目标日期、窗口、模式、状态、总数、有效数、覆盖率和父重试批次 |
| `market_sync_item` | `(batch_id, item_kind, item_key)` | 股票行情或指标日期等同步对象的状态、计数、指纹和脱敏错误摘要 |
行情和每日指标按日期建查询索引;数据库约束保证幂等。金额、价格和比率采用能稳定往返 Tushare 有效精度的 `NUMERIC`,避免把二进制浮点文本差异带入指纹或相等判断。迁移由 Alembic 独立执行,FastAPI 启动时不自动建表。
### 3.3 CSV 布局
CSV 根目录由配置指定,候选布局如下:
```text
market-data/
├── bars/<ts_code>.csv
├── daily-basic/<YYYY>/<YYYYMMDD>.csv
└── stock-basic/current.csv
```
- 每股 bar 文件始终是当前六年 qfq 正式快照。
- `daily-basic` 按交易日落地,初始化可逐日恢复,日常只发布目标交易日;成功批次清理六年窗口之前的日期文件。
- 临时文件与正式文件必须位于同一文件系统,使用 flush、fsync 和 `os.replace` 完成原子晋升。
- CSV 表头、列顺序、日期格式、空值和小数序列化固定,避免运行环境差异改变指纹。
## 4. 同步数据流
### 4.1 批次准备
1. 获取交易日历并确定目标交易日和六年窗口。
2. 刷新 `stock_basic(list_status="L")`,按代码段保留沪深 A 股并排除 ST/风险警示名称;在 PostgreSQL 中把当前目标股票集合标记为 active,其他历史记录不参与本批次。
3. 创建 `running` 批次。普通运行覆盖全部目标股票;重试运行读取父批次失败的股票或指标日期并创建新批次,不修改原批次审计记录。
### 4.2 每日指标
- 空库初始化按开市日调用 `daily_basic(trade_date=...)` 回补六年,并以“一个交易日”作为文件发布和数据库事务边界。
- 日常运行只请求目标交易日;重复运行通过唯一键和差异更新保持幂等。
- 每个日期先过滤到当前目标股票集合,再写临时 CSV 并校验代码唯一性、日期一致性和必要字段;随后在一个事务中通过 `COPY` staging + 集合式 upsert 写库,提交后发布 CSV。
- 日期写入失败时保留旧库和旧正式文件,并作为批次失败数据源参与覆盖率计算。
### 4.3 每股 qfq 日线
对每个目标股票执行:
1. 调用 `pro_bar(adj="qfq", start_date=window_start, end_date=target_trade_date)`,规范化、按日期升序、检查唯一键与 OHLC 合法性,写入临时 CSV。
2. 若不存在旧正式 CSV,走完整窗口写入。
3. 若旧快照存在,取新旧共同日期边界形成完整重叠区间;分别序列化该区间内的固定行情列和日期集合并计算 SHA-256。缺行、增行或任一字段变化都会改变指纹。
4. 指纹相同时,仅把新快照中数据库窗口上界之后的日期批量插入,并在同一事务中清理 `< window_start` 的旧行。
5. 指纹不同时,把新窗口 `COPY` 到临时 staging 表,集合式 upsert 改变行,并删除该股票在窗口内已不再存在于新快照的行以及 `< window_start` 的旧行。
6. 数据库提交后原子替换正式 CSV;异常时回滚该股票事务、删除临时文件、保留旧正式 CSV,并记录失败项。
指纹只覆盖 `trade_date` 和固定行情字段,不包含批次 ID、抓取时间、当前股票名称或每日估值数据。SHA-256 的用途是变化检测而非安全认证。
### 4.4 批次收敛与覆盖率
- 每只股票成功或失败独立提交,不因少数失败回滚已成功股票。
- 目标股票只有同时存在目标交易日的 bar 与策略所需 `daily_basic` 才计入有效集合。
- `coverage = valid_stock_count / target_stock_count`;批次记录目标数、有效数、覆盖率和失败列表。
- 全部目标对象成功为 `success`;至少一个对象成功且存在失败为 `partial_success`;没有可用成功结果或批次准备失败为 `failed`。
- 本任务只返回 `strategy_eligible = coverage >= configured_threshold` 等契约,不调用选股策略。后续人工强制运行由策略任务记录“不完整”标记。
## 5. PostgreSQL 写入与事务
- 运行时使用 Psycopg 3,schema 迁移使用 Alembic + SQLAlchemy metadata;业务同步不引入 ORM 实体生命周期。
- 每次批量写入创建连接级或事务级临时 staging 表,通过 Psycopg `COPY FROM STDIN` 装载,再执行 `INSERT ... ON CONFLICT DO UPDATE ... WHERE target IS DISTINCT FROM excluded`。
- bar 事务边界是单只股票;daily basic 事务边界是单个交易日;批次与 item 状态使用短事务记录,避免一个长事务覆盖整个市场。
- 同一环境只允许一个市场同步运行。使用 PostgreSQL advisory lock 快速拒绝重叠执行,避免 cron 重叠造成 CSV 发布竞争;退出码区分成功、部分成功/覆盖率不足和基础设施失败。
- 连接凭据、Tushare token 和原始异常中的秘密不写入日志或批次错误字段。
## 6. Compose 与外部调度
- Compose 增加带健康检查和命名数据卷的 PostgreSQL 服务;`ZHIXING_DATABASE_URL` 可覆盖默认内部连接,便于未来改用外部托管数据库。
- 增加复用 production 后端镜像的 `migrate` 与 `market-sync` 一次性服务,放入 `jobs` profile;正常 `docker compose up -d` 不会常驻启动它们。
- `market-sync` 挂载持久 CSV 卷并执行 CLI。部署文档先运行 migration,再给出宿主机 cron 的 `docker compose --profile jobs run --rm market-sync` 示例。
- cron 是宿主机自带能力,不需要额外 Docker 应用;任务本身通过交易日历、唯一约束和 advisory lock 应对节假日、重复或重叠触发。
## 7. 错误处理、可观测性与恢复
- Tushare 网络/限流错误采用有上限的指数退避和抖动;并发数、重试次数、请求间隔均由配置控制,默认保守,避免把旧项目固定八线程作为隐式契约。
- 校验错误、供应商错误和数据库错误映射为明确 item 失败类型;日志使用 batch ID、`ts_code`、阶段和计数,不记录 token 或整批原始响应。
- `partial_success` 后可按父 batch ID 重试失败股票和失败指标日期;已成功对象不会重新拉取,除非启动新的普通批次。
- 正式 CSV 是最后一次数据库成功提交后的可恢复快照。若 CSV 晋升在数据库提交后极少数失败,item 标记为发布失败;重试会以数据库幂等写入后再次发布,不能宣称整体成功。
- 回滚代码时保留 PostgreSQL volume 和 CSV volume;数据库 schema 回退仅使用经过审查的 Alembic downgrade,任何数据删除前先备份。
## 8. 重要权衡与延期项
- 每股请求六年与增量请求消耗相同调用次数,换取 qfq 历史修订检测;本地约 6.86M 行 CSV 的 SHA-1 全量测试约 1.87 秒,指纹不是主要性能瓶颈,网络调用和极端全量修复才是重点。
- 当前股票池存在幸存者偏差,且 qfq 使用最新修订值;这是已接受的近期走势分析边界,不适合声称长期历史回测可复现。
- 暂不实现数据库分区、任务队列、前端页面、策略运行器和跨供应商抽象。六年数据量可先由普通索引表承载,获得真实运行指标后再决定是否分区。
## 9. 兼容与迁移
- 不直接导入旧项目把最新总市值复制到历史行的 CSV;行情可重新从 Tushare 拉取,daily basic 按交易日回补,避免继承前视偏差。
- 新 CSV 格式是新的落地契约,旧项目文件只作人工核对,不作为自动迁移输入。
- 当前 FastAPI HTTP 契约不变;server 即使不触发同步也可继续提供现有健康检查。