Merge pull request 'Develop' (#2) from develop into main
Deploy Production / deploy (push) Successful in 18s
Deploy Production / deploy (push) Successful in 18s
Reviewed-on: sakibcc/zhixing-system#2
This commit was merged in pull request #2.
This commit is contained in:
@@ -0,0 +1,5 @@
|
|||||||
|
{"file":".trellis/spec/backend/index.md","reason":"检查实现是否符合后端总入口的开发前检查与质量门禁。"}
|
||||||
|
{"file":".trellis/spec/backend/directory-structure.md","reason":"检查 market_data 分层和导入方向是否符合 bounded context 约束。"}
|
||||||
|
{"file":".trellis/spec/backend/configuration-and-runtime.md","reason":"检查配置、应用生命周期与容器运行边界。"}
|
||||||
|
{"file":".trellis/spec/backend/quality-guidelines.md","reason":"执行并核对 Ruff、Pyright、pytest 及禁止模式。"}
|
||||||
|
{"file":".trellis/tasks/08-05-tushare-postgres-sync/research/current-apis.md","reason":"核对 COPY、事务、迁移和 Compose profile 实现是否符合当前官方接口。"}
|
||||||
@@ -0,0 +1,134 @@
|
|||||||
|
# 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 即使不触发同步也可继续提供现有健康检查。
|
||||||
@@ -0,0 +1,5 @@
|
|||||||
|
{"file":".trellis/spec/backend/index.md","reason":"实现前掌握后端 bounded context、运行时与质量入口。"}
|
||||||
|
{"file":".trellis/spec/backend/directory-structure.md","reason":"按 domain/application/infrastructure/presentation 边界创建 market_data 模块。"}
|
||||||
|
{"file":".trellis/spec/backend/configuration-and-runtime.md","reason":"数据库、Tushare、CSV 与同步参数必须进入统一 Settings 配置入口。"}
|
||||||
|
{"file":".trellis/spec/backend/quality-guidelines.md","reason":"遵循 Python 3.12 strict typing、Ruff、Pyright 与 pytest 约束。"}
|
||||||
|
{"file":".trellis/tasks/08-05-tushare-postgres-sync/research/current-apis.md","reason":"使用已核对的 Tushare、Psycopg COPY、Alembic 与 Compose Job 接口。"}
|
||||||
@@ -0,0 +1,129 @@
|
|||||||
|
# Tushare PostgreSQL 同步实施计划
|
||||||
|
|
||||||
|
## 实施原则
|
||||||
|
|
||||||
|
- 按以下顺序增量交付,每一步先写对应行为测试,再补最小实现。
|
||||||
|
- 不修改旧项目;不导入其带有最新市值回填的历史 CSV。
|
||||||
|
- 不在本任务中实现选股策略、HTTP 管理接口或宿主机 crontab 写入。
|
||||||
|
- schema、CSV 契约和批次状态一旦被后续步骤使用,变更时必须同步更新迁移、测试和设计文档。
|
||||||
|
|
||||||
|
## 1. 依赖、配置与迁移骨架
|
||||||
|
|
||||||
|
- 在 `zhixing-server/pyproject.toml` 增加 Tushare、Psycopg 3、SQLAlchemy 和 Alembic 依赖并更新 `uv.lock`。
|
||||||
|
- 扩展 `bootstrap.config.Settings`:数据库 URL、Tushare token、CSV 根目录、覆盖率阈值、并发/限流/重试参数;保持 `ZHIXING_` 前缀且不暴露秘密。
|
||||||
|
- 初始化 Alembic 配置和 metadata,创建股票、日线、每日指标、同步 batch/item 表及唯一约束、索引。
|
||||||
|
- 增加 migration 集成测试,验证空库 upgrade 到 head、唯一约束和必要索引。
|
||||||
|
|
||||||
|
验证:
|
||||||
|
|
||||||
|
```bash
|
||||||
|
cd zhixing-server
|
||||||
|
uv sync
|
||||||
|
uv run alembic upgrade head
|
||||||
|
uv run pytest tests/integration/test_market_data_migrations.py
|
||||||
|
```
|
||||||
|
|
||||||
|
回滚点:本阶段只建立 schema;若迁移失败,修正 migration 后重建测试库,不对生产卷执行手工 DDL。
|
||||||
|
|
||||||
|
## 2. 领域模型、规范化与端口
|
||||||
|
|
||||||
|
- 创建 `modules/market_data` 四层目录和 bounded-context README。
|
||||||
|
- 定义 `ts_code`、目标交易日、六年窗口、bar、daily basic、batch/item 状态和值对象。
|
||||||
|
- 定义 Tushare 来源、快照存储、市场数据仓储、批次仓储和同步锁端口。
|
||||||
|
- 实现固定列、日期、小数、空值规范化,完整重叠区间 SHA-256 计算和变化分类。
|
||||||
|
- 先用纯领域测试覆盖:初次同步、仅新增日期、历史值变化、重叠区间缺行、输入顺序变化、空值/小数稳定性和窗口边界。
|
||||||
|
|
||||||
|
验证:
|
||||||
|
|
||||||
|
```bash
|
||||||
|
cd zhixing-server
|
||||||
|
uv run pytest tests/unit/market_data/test_fingerprint.py tests/unit/market_data/test_window.py
|
||||||
|
```
|
||||||
|
|
||||||
|
## 3. Tushare 与 CSV 适配器
|
||||||
|
|
||||||
|
- 封装 `stock_basic`、`trade_cal`、`pro_bar(adj="qfq")` 和 `daily_basic(trade_date=...)`,把 DataFrame 转为领域记录。
|
||||||
|
- 实现当前沪深非 ST 股票池过滤、请求退避/抖动和可配置并发限制。
|
||||||
|
- 实现 bar 按股票、daily basic 按日期、stock basic 当前版的临时 CSV 写入、校验、指纹读取、原子晋升和过期文件清理。
|
||||||
|
- 通过 fake Tushare 客户端与临时目录测试,不让普通单元测试访问真实网络。
|
||||||
|
|
||||||
|
验证:
|
||||||
|
|
||||||
|
```bash
|
||||||
|
cd zhixing-server
|
||||||
|
uv run pytest tests/unit/market_data/test_tushare_adapter.py tests/unit/market_data/test_csv_snapshot_store.py
|
||||||
|
```
|
||||||
|
|
||||||
|
回滚点:正式 CSV 只在数据库提交后晋升;适配器测试必须证明异常路径保留旧文件。
|
||||||
|
|
||||||
|
## 4. PostgreSQL 批量仓储
|
||||||
|
|
||||||
|
- 使用 Psycopg 连接/事务实现股票主数据、bar、daily basic、batch/item 和 advisory lock 适配器。
|
||||||
|
- 对 bar 和 daily basic 使用临时 staging + `COPY` + 集合式 upsert;通过 `IS DISTINCT FROM` 避免无意义 update。
|
||||||
|
- 实现 insert-only、新窗口完整修复、窗口内缺失删除和六年前数据清理。
|
||||||
|
- 集成测试覆盖幂等、更新计数、单股票回滚、日期事务回滚、并发锁以及数据库最终与快照一致。
|
||||||
|
|
||||||
|
验证:
|
||||||
|
|
||||||
|
```bash
|
||||||
|
cd zhixing-server
|
||||||
|
uv run pytest tests/integration/test_market_data_repository.py
|
||||||
|
```
|
||||||
|
|
||||||
|
## 5. 同步用例与失败重试
|
||||||
|
|
||||||
|
- 实现 `SyncMarketData.execute()`:解析目标交易日、刷新股票池、创建批次、回补/更新 daily basic、逐股票同步 bar、收敛状态和覆盖率。
|
||||||
|
- 初始化模式只回补缺失开市日;日常模式只获取目标日 daily basic,但每股始终拉六年 qfq。
|
||||||
|
- 实现父 batch 的失败股票和失败指标日期重试;成功项不回滚、不重复抓取。
|
||||||
|
- 用 fake 端口编排测试覆盖 `success`、`partial_success`、`failed`、99% 阈值、旧数据不得计入有效集合、CSV 发布失败和重试恢复。
|
||||||
|
|
||||||
|
验证:
|
||||||
|
|
||||||
|
```bash
|
||||||
|
cd zhixing-server
|
||||||
|
uv run pytest tests/unit/market_data/test_sync_market_data.py
|
||||||
|
```
|
||||||
|
|
||||||
|
## 6. CLI、Compose 与运维文档
|
||||||
|
|
||||||
|
- 增加 console script/CLI,支持普通同步、显式目标日期、初始化和按 batch 重试;输出批次摘要并返回稳定退出码。
|
||||||
|
- 在 dev/prod Compose 中增加 PostgreSQL、CSV 持久卷、`migrate` 和 `market-sync` profile Job,复用后端镜像和统一环境配置。
|
||||||
|
- 更新 `.env.example` 与部署文档,提供 migration、手工运行、失败重试和宿主机 cron 示例;不自动写入 crontab。
|
||||||
|
- 验证正常 `up` 不启动 Job,显式 `run --rm` 能连接数据库与 CSV 卷。
|
||||||
|
|
||||||
|
验证:
|
||||||
|
|
||||||
|
```bash
|
||||||
|
docker compose -f docker-compose.prod.yml config
|
||||||
|
docker compose -f docker-compose.dev.yml config
|
||||||
|
docker compose -f docker-compose.prod.yml --profile jobs config
|
||||||
|
```
|
||||||
|
|
||||||
|
## 7. 全链路验收与质量门禁
|
||||||
|
|
||||||
|
- 使用受控 fixture 或测试 Tushare 适配器完成空库初始化、无变化重跑、仅新增日期、历史 qfq 修订、单股票失败和定向重试场景。
|
||||||
|
- 核对 batch/item 计数、覆盖率、数据库内容、正式 CSV、六年清理和秘密脱敏。
|
||||||
|
- 若运行真实 Tushare smoke test,必须由显式环境开关启用,限制少量股票且不把 token 写入测试输出。
|
||||||
|
- 运行后端和根目录质量门禁;修复所有 Ruff、Pyright、pytest 与 Compose config 问题。
|
||||||
|
|
||||||
|
验证:
|
||||||
|
|
||||||
|
```bash
|
||||||
|
cd zhixing-server
|
||||||
|
uv run ruff format --check .
|
||||||
|
uv run ruff check .
|
||||||
|
uv run pyright
|
||||||
|
uv run pytest
|
||||||
|
cd ..
|
||||||
|
./dev.sh check
|
||||||
|
./dev.sh test
|
||||||
|
git diff --check
|
||||||
|
```
|
||||||
|
|
||||||
|
## 开始实施前检查
|
||||||
|
|
||||||
|
- [ ] 用户已在最终规划摘要之后明确批准实施。
|
||||||
|
- [ ] `prd.md`、`design.md`、`implement.md` 内容一致且无阻塞问题。
|
||||||
|
- [ ] `implement.jsonl` 与 `check.jsonl` 各自至少包含一条真实规格/研究上下文。
|
||||||
|
- [ ] `task.py validate tushare-postgres-sync` 通过。
|
||||||
|
- [ ] 当前工作区已有改动已识别,实施不覆盖用户或其他任务的文件。
|
||||||
@@ -0,0 +1,102 @@
|
|||||||
|
# 迁移 Tushare PostgreSQL 同步
|
||||||
|
|
||||||
|
## Goal
|
||||||
|
|
||||||
|
把旧项目中面向当前上市股票的 Tushare 日线获取能力迁移到 `zhixing-server`,形成可由外部调度器重复执行的一次性同步任务。同步结果以 PostgreSQL 为策略查询事实源,同时保留按股票拆分的六年 qfq CSV 落地快照,以可靠识别和修复历史复权数据变化。
|
||||||
|
|
||||||
|
## User Value
|
||||||
|
|
||||||
|
- 每个交易日收盘后可以稳定更新当前股票池的日线、股票主数据和每日估值/交易指标。
|
||||||
|
- 后续选股策略读取结构化、可查询且有覆盖率保障的市场数据,不再直接依赖散落 CSV。
|
||||||
|
- Tushare 修订历史 qfq 时,可以只修复受影响股票,而不会默默保留过期历史价格。
|
||||||
|
|
||||||
|
## Background and Confirmed Facts
|
||||||
|
|
||||||
|
- 当前 `zhixing-server` 只有 FastAPI 基础工程,尚无数据库连接、迁移、任务队列或市场数据 bounded context。
|
||||||
|
- 旧项目按股票调用 `pro_bar(adj="qfq")` 并写 CSV,但初始化流程总是全量覆盖,没有 PostgreSQL 行情存储和可审计同步批次。
|
||||||
|
- 旧项目把最新总市值复制到所有历史 K 线行,不满足按交易日分析估值条件的需求;新实现必须把每日估值/交易指标按交易日保存。
|
||||||
|
- 已接受 [PostgreSQL 主存储决策](../../../docs/adr/0003-postgresql-as-market-data-store.md) 和 [六年快照同步决策](../../../docs/adr/0004-tushare-six-year-snapshot-sync.md)。
|
||||||
|
|
||||||
|
## Requirements
|
||||||
|
|
||||||
|
### R1. 数据范围
|
||||||
|
|
||||||
|
- 同步当前上市的沪深非 ST A 股股票主数据;排除 ST/风险警示股和北交所股票。
|
||||||
|
- 每只股票同步最近六年的日线行情,价格口径只保留 Tushare `qfq`。
|
||||||
|
- 同步按交易日归属的估值与交易指标,至少满足后续市值、流动性和换手率筛选。
|
||||||
|
- 只支持日线收盘数据,不包含分钟级、盘中实时行情或长期无幸存者偏差回测数据。
|
||||||
|
|
||||||
|
### R2. 存储职责
|
||||||
|
|
||||||
|
- PostgreSQL 是运行时和策略查询的市场数据事实源。
|
||||||
|
- CSV 是 Tushare 落地快照、导出和恢复介质,不承担策略在线查询。
|
||||||
|
- 日线数据以 Tushare `(ts_code, trade_date)` 唯一标识;重复同步必须幂等。
|
||||||
|
- PostgreSQL 与正式 CSV 均按目标交易日滚动保留最近六年;成功同步新窗口时清理窗口起点之前的行情、每日指标和日期快照。
|
||||||
|
- 某个对象同步失败时不得因本次运行提前删除其原有数据或替换其正式 CSV。
|
||||||
|
|
||||||
|
### R3. 六年快照与修订检测
|
||||||
|
|
||||||
|
- 每次同步按股票重新请求最近六年 qfq 日线,而不是只请求最新日期。
|
||||||
|
- 新旧 CSV 必须在规范化后的完整历史重叠区间计算并比较指纹。
|
||||||
|
- 指纹相同时,只向 PostgreSQL 插入新交易日。
|
||||||
|
- 指纹不同时,只对发生变化的股票执行最近六年完整 upsert,不触发全市场回写。
|
||||||
|
- 指纹字段不得包含同步时间、批次 ID 或会使历史行每天变化的当前快照字段。
|
||||||
|
|
||||||
|
### R4. 原子性与恢复
|
||||||
|
|
||||||
|
- 单只股票按“临时 CSV → 校验和指纹比较 → PostgreSQL 事务提交 → 原子替换正式 CSV”的顺序处理。
|
||||||
|
- PostgreSQL 失败时不得替换该股票的正式 CSV。
|
||||||
|
- PostgreSQL 批量修复必须使用批量装载和集合式 upsert,不得通过 ORM 逐行更新六年数据。
|
||||||
|
|
||||||
|
### R5. 批次、部分成功与重试
|
||||||
|
|
||||||
|
- 单只股票是事务和重试边界。
|
||||||
|
- 同步批次支持 `success`、`partial_success` 和 `failed` 状态,并记录每只股票行情及每日指标日期的插入、更新、未变化或失败结果。
|
||||||
|
- 部分失败不得回滚成功对象;后续可以只重试失败股票或失败指标日期。
|
||||||
|
|
||||||
|
### R6. 数据新鲜度与后续选股
|
||||||
|
|
||||||
|
- 有效选股股票必须同时具备目标交易日的行情和所需估值数据,旧日期数据不得冒充当天数据。
|
||||||
|
- 自动选股最低数据覆盖率默认 `99%`,并可配置。
|
||||||
|
- 覆盖率不足时不自动运行策略;人工强制运行必须把结果标记为不完整。
|
||||||
|
- 本任务只建立覆盖率与新鲜度契约,不迁移或实现具体选股策略。
|
||||||
|
|
||||||
|
### R7. 运行边界
|
||||||
|
|
||||||
|
- 同步由外部调度器触发,FastAPI 进程内部不运行定时器。
|
||||||
|
- 项目在 Docker Compose 中定义复用后端镜像的一次性同步 Job,并通过 CLI 执行可重复的一次性同步用例。
|
||||||
|
- 生产宿主机 cron 使用 `docker compose run --rm` 定时启动 Job;本任务提供配置示例,但不直接修改宿主机 crontab。
|
||||||
|
- Tushare 凭据和 PostgreSQL 连接信息通过项目统一配置入口注入,不写入代码、日志或响应。
|
||||||
|
|
||||||
|
## Technical Notes
|
||||||
|
|
||||||
|
- Tushare `pro_bar(adj="qfq")` 内部组合日线与复权因子;六年范围的每股请求适合作为 qfq 修订检测来源。
|
||||||
|
- Tushare `daily_basic` 支持按 `trade_date` 获取全市场、单次最多 6000 行,更新时间为交易日 15:00~17:00;首次初始化应按交易日回补,日常只同步目标交易日,不需要每天重拉六年估值指标。
|
||||||
|
- PostgreSQL 批量修复计划使用 Psycopg 3 `COPY FROM STDIN` 加 staging 表和集合式 upsert;数据库 schema 使用独立迁移管理,不在 FastAPI 启动时自动建表。
|
||||||
|
- Docker Compose Job 使用 profile 与 `docker compose run --rm`,宿主机 cron 仅负责触发。
|
||||||
|
|
||||||
|
## Acceptance Criteria
|
||||||
|
|
||||||
|
- [ ] 可在空 PostgreSQL 中完成当前目标股票池的六年 qfq 日线、股票主数据和每日估值/交易指标初始化。
|
||||||
|
- [ ] 目标股票池只包含当前上市的沪深非 ST A 股,不包含 ST/风险警示股和北交所股票。
|
||||||
|
- [ ] 对相同 Tushare 返回重复执行同步,不产生重复行或无意义历史更新。
|
||||||
|
- [ ] 只有新交易日时,仅新增对应日期数据。
|
||||||
|
- [ ] 任意历史重叠行变化时,只对该股票执行六年修复,PostgreSQL 最终与新 CSV 一致。
|
||||||
|
- [ ] PostgreSQL 写入失败时,该股票正式 CSV 保持旧版本;重试后可以幂等完成。
|
||||||
|
- [ ] 少数股票失败时批次为 `partial_success`,成功股票保留结果,失败股票可单独重试。
|
||||||
|
- [ ] 同步结果能够报告目标数、有效数、覆盖率、插入数、更新数、未变化数和失败列表。
|
||||||
|
- [ ] 覆盖率低于默认 `99%` 时不会自动触发选股;人工强制运行具有显式不完整标记。
|
||||||
|
- [ ] Compose 一次性 Job 可以复用生产后端镜像、数据库连接和 CSV 数据卷,正常 `docker compose up -d` 不会常驻启动该 Job。
|
||||||
|
- [ ] 项目文档提供宿主机 cron 调用 `docker compose run --rm` 的示例,周末、节假日和重复触发保持安全幂等。
|
||||||
|
- [ ] 成功同步后 PostgreSQL 与正式 CSV 不包含目标六年窗口之前的数据;失败对象仍保留上一次成功发布的数据。
|
||||||
|
- [ ] 单元测试覆盖指纹、差异判定、幂等写入、失败恢复和覆盖率计算;PostgreSQL 集成测试覆盖迁移、唯一约束和批量 upsert。
|
||||||
|
- [ ] 后端 Ruff、Pyright 和 pytest 质量门禁通过。
|
||||||
|
|
||||||
|
## Out of Scope
|
||||||
|
|
||||||
|
- 分钟级、盘中实时行情和 WebSocket 数据。
|
||||||
|
- 未复权、后复权或复权因子长期保存。
|
||||||
|
- 退市股票历史成员资格和无幸存者偏差的长期回测股票池。
|
||||||
|
- 具体选股策略、信号持久化、图表和前端管理页面。
|
||||||
|
- FastAPI 进程内定时器或本任务直接配置生产调度平台。
|
||||||
|
- 六年以上历史的长期保存、退市样本回补和长期回测支持。
|
||||||
@@ -0,0 +1,20 @@
|
|||||||
|
# 当前接口与工具核对
|
||||||
|
|
||||||
|
## Tushare
|
||||||
|
|
||||||
|
- `pro_bar(adj="qfq")` 会请求日线与复权因子,并以最新因子归一化价格;因此新增除权事件可能改变请求区间内的历史 qfq。[Tushare `pro_bar` 源码](https://github.com/waditu/tushare/blob/master/tushare/pro/data_pro.py)
|
||||||
|
- `daily_basic` 支持按股票或交易日请求,单次最多返回 6000 行,交易日 15:00~17:00 更新;按交易日循环适合六年初始化,日常同步只需请求目标交易日。[Tushare 每日指标](https://tushare.pro/document/2?doc_id=32)
|
||||||
|
|
||||||
|
## PostgreSQL 与 Psycopg
|
||||||
|
|
||||||
|
- PostgreSQL 官方建议批量装载使用 `COPY`,其性能显著优于连续单行 `INSERT`。[PostgreSQL Populating a Database](https://www.postgresql.org/docs/current/populate.html)
|
||||||
|
- Psycopg 3 通过 `Cursor.copy("COPY ... FROM STDIN")` 和上下文管理器写入批量数据;异常会遵循事务回滚语义。[Psycopg COPY](https://www.psycopg.org/psycopg3/docs/basic/copy.html)
|
||||||
|
- Psycopg 3 支持 `with conn.transaction()`,异常时回滚,适合作为单股票同步事务边界。[Psycopg transactions](https://www.psycopg.org/psycopg3/docs/api/connections.html)
|
||||||
|
|
||||||
|
## 数据库迁移
|
||||||
|
|
||||||
|
- Alembic 可以通过 `target_metadata` 管理 SQLAlchemy schema,并可检查数据库 revision 是否到达 head;迁移应作为独立部署步骤执行。[Alembic autogenerate](https://alembic.sqlalchemy.org/en/latest/autogenerate.html)
|
||||||
|
|
||||||
|
## Docker Compose
|
||||||
|
|
||||||
|
- Docker 官方示例使用 profile 定义不会随正常 `up` 启动的一次性服务,并通过 `docker compose run --rm` 执行临时任务。[Docker Compose profiles](https://github.com/docker/docs/blob/main/content/manuals/compose/how-tos/profiles.md)
|
||||||
@@ -0,0 +1,34 @@
|
|||||||
|
{
|
||||||
|
"id": "tushare-postgres-sync",
|
||||||
|
"name": "tushare-postgres-sync",
|
||||||
|
"title": "迁移 Tushare PostgreSQL 同步",
|
||||||
|
"description": "将旧项目的 Tushare 日线同步迁移到 PostgreSQL,保留六年 qfq CSV 落地快照并支持历史修订检测。",
|
||||||
|
"status": "completed",
|
||||||
|
"dev_type": "backend",
|
||||||
|
"scope": "zhixing-server market data synchronization",
|
||||||
|
"package": "zhixing-server",
|
||||||
|
"priority": "P2",
|
||||||
|
"creator": "yuxuanhui",
|
||||||
|
"assignee": "yuxuanhui",
|
||||||
|
"createdAt": "2026-08-05",
|
||||||
|
"completedAt": "2026-08-06",
|
||||||
|
"branch": null,
|
||||||
|
"base_branch": "main",
|
||||||
|
"worktree_path": null,
|
||||||
|
"commit": null,
|
||||||
|
"pr_url": null,
|
||||||
|
"subtasks": [],
|
||||||
|
"children": [],
|
||||||
|
"parent": null,
|
||||||
|
"relatedFiles": [
|
||||||
|
"CONTEXT.md",
|
||||||
|
"docs/adr/0003-postgresql-as-market-data-store.md",
|
||||||
|
"docs/adr/0004-tushare-six-year-snapshot-sync.md",
|
||||||
|
"zhixing-server/src/zhixing_server/modules/",
|
||||||
|
"zhixing-server/pyproject.toml",
|
||||||
|
"docker-compose.dev.yml",
|
||||||
|
"docker-compose.prod.yml"
|
||||||
|
],
|
||||||
|
"notes": "",
|
||||||
|
"meta": {}
|
||||||
|
}
|
||||||
@@ -8,8 +8,8 @@
|
|||||||
|
|
||||||
<!-- @@@auto:current-status -->
|
<!-- @@@auto:current-status -->
|
||||||
- **Active File**: `journal-1.md`
|
- **Active File**: `journal-1.md`
|
||||||
- **Total Sessions**: 0
|
- **Total Sessions**: 1
|
||||||
- **Last Active**: -
|
- **Last Active**: 2026-08-06
|
||||||
<!-- @@@/auto:current-status -->
|
<!-- @@@/auto:current-status -->
|
||||||
|
|
||||||
---
|
---
|
||||||
@@ -19,7 +19,7 @@
|
|||||||
<!-- @@@auto:active-documents -->
|
<!-- @@@auto:active-documents -->
|
||||||
| File | Lines | Status |
|
| File | Lines | Status |
|
||||||
|------|-------|--------|
|
|------|-------|--------|
|
||||||
| `journal-1.md` | ~0 | Active |
|
| `journal-1.md` | ~29 | Active |
|
||||||
<!-- @@@/auto:active-documents -->
|
<!-- @@@/auto:active-documents -->
|
||||||
|
|
||||||
---
|
---
|
||||||
@@ -29,6 +29,7 @@
|
|||||||
<!-- @@@auto:session-history -->
|
<!-- @@@auto:session-history -->
|
||||||
| # | Date | Title | Commits | Branch |
|
| # | Date | Title | Commits | Branch |
|
||||||
|---|------|-------|---------|--------|
|
|---|------|-------|---------|--------|
|
||||||
|
| 1 | 2026-08-06 | 完成 TUSHARE PostgreSQL 同步任务 | `039a81f`, `11e8728` | `develop` |
|
||||||
<!-- @@@/auto:session-history -->
|
<!-- @@@/auto:session-history -->
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|||||||
@@ -5,3 +5,25 @@
|
|||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
|
## Session 1: 完成 TUSHARE PostgreSQL 同步任务
|
||||||
|
|
||||||
|
**Date**: 2026-08-06
|
||||||
|
**Task**: 完成 TUSHARE PostgreSQL 同步任务
|
||||||
|
**Branch**: `develop`
|
||||||
|
|
||||||
|
### Summary
|
||||||
|
|
||||||
|
完成 PostgreSQL 市场数据同步、TUSHARE 真实单股票验证、开发与生产配置、规格文档和回归测试;已归档当前同步任务。
|
||||||
|
|
||||||
|
### Git Commits
|
||||||
|
|
||||||
|
| Hash | Message |
|
||||||
|
|------|---------|
|
||||||
|
| `039a81f` | (see git log) |
|
||||||
|
| `11e8728` | (see git log) |
|
||||||
|
|
||||||
|
### Status
|
||||||
|
|
||||||
|
[OK] **Completed**
|
||||||
|
|||||||
+28
@@ -0,0 +1,28 @@
|
|||||||
|
# 知行量化研究上下文
|
||||||
|
|
||||||
|
本上下文描述 A 股市场数据同步与选股策略分析中的业务语言,重点保证历史结果可以被可靠复现。
|
||||||
|
|
||||||
|
## 分析语义
|
||||||
|
|
||||||
|
**历史可复现分析**:给定一个历史交易日,在当前目标股票池范围内重建该日的行情与估值状态;价格采用数据库中当前最新修订的前复权行情,不要求还原某次同步时的价格版本。
|
||||||
|
_Avoid_: 用当前市值回填历史、把当前目标股票池当作无幸存者偏差的长期回测宇宙
|
||||||
|
|
||||||
|
## 市场数据
|
||||||
|
|
||||||
|
**股票主数据**:描述证券标识及其在市场中的基础状态,例如代码、名称、交易所和上市状态。
|
||||||
|
|
||||||
|
**当前目标股票池**:同步时仍处于上市状态的沪深 A 股集合,排除 ST/风险警示股和北交所股票,也不包含退市股票的长期历史成员资格。
|
||||||
|
|
||||||
|
**有效选股股票池**:当前目标股票池中,行情和所需估值数据均已同步到目标交易日的股票集合;数据仍停留在更早交易日的股票不参与当天选股。
|
||||||
|
|
||||||
|
**日线行情**:以证券和交易日为粒度记录的开盘价、最高价、最低价、收盘价和成交量等交易结果。
|
||||||
|
|
||||||
|
**前复权行情**:按 `qfq` 口径整理的日线价格序列,作为当前选股策略分析的唯一价格口径。
|
||||||
|
|
||||||
|
**估值快照**:在某个交易日可见的证券估值与交易活跃度指标,例如总市值、流通市值和换手率。
|
||||||
|
|
||||||
|
**行情落地快照**:一次同步从数据源取得的完整行情集合,用于核对和传递数据,不承担策略查询职责。
|
||||||
|
|
||||||
|
**同步批次**:一次面向当前目标股票池的市场数据同步运行;批次可以全部成功、部分成功或失败,并保留每只股票的处理结果。
|
||||||
|
|
||||||
|
**数据覆盖率**:目标交易日内,有效选股股票数占当前目标股票池目标数的比例,用于判断选股结果是否具备足够完整性。
|
||||||
@@ -51,11 +51,28 @@ services:
|
|||||||
context: ./zhixing-server
|
context: ./zhixing-server
|
||||||
target: production
|
target: production
|
||||||
command: ["alembic", "upgrade", "head"]
|
command: ["alembic", "upgrade", "head"]
|
||||||
|
depends_on:
|
||||||
|
market-data-init:
|
||||||
|
condition: service_completed_successfully
|
||||||
environment:
|
environment:
|
||||||
ZHIXING_DATABASE_URL: ${ZHIXING_DATABASE_URL:?Set ZHIXING_DATABASE_URL to the 1Panel PostgreSQL URL}
|
ZHIXING_DATABASE_URL: ${ZHIXING_DATABASE_URL:?Set ZHIXING_DATABASE_URL to the 1Panel PostgreSQL URL}
|
||||||
networks:
|
networks:
|
||||||
- 1panel-network
|
- 1panel-network
|
||||||
|
|
||||||
|
# Named volumes are initialized as root; normalize ownership before app jobs run.
|
||||||
|
market-data-init:
|
||||||
|
profiles: ["jobs"]
|
||||||
|
build:
|
||||||
|
context: ./zhixing-server
|
||||||
|
target: production
|
||||||
|
user: "0:0"
|
||||||
|
command:
|
||||||
|
- sh
|
||||||
|
- -c
|
||||||
|
- mkdir -p /app/data/market-data && chown -R 10001:10001 /app/data/market-data
|
||||||
|
volumes:
|
||||||
|
- market-data:/app/data/market-data
|
||||||
|
|
||||||
market-sync:
|
market-sync:
|
||||||
profiles: ["jobs"]
|
profiles: ["jobs"]
|
||||||
build:
|
build:
|
||||||
|
|||||||
@@ -0,0 +1,25 @@
|
|||||||
|
---
|
||||||
|
status: accepted
|
||||||
|
---
|
||||||
|
|
||||||
|
# 市场数据以 PostgreSQL 为主存储
|
||||||
|
|
||||||
|
为支持历史可复现分析、按交易日查询、增量同步和跨标的策略分析,标准化的股票池、日线行情与估值数据统一落入 PostgreSQL。CSV 不作为运行时事实源,仅作为 Tushare 落地快照、导出和故障恢复介质;策略层通过仓储接口读取 DataFrame,以隔离分析逻辑与存储技术。
|
||||||
|
|
||||||
|
第一阶段同步范围包含股票主数据、日线行情和每日估值/交易指标;只迁移 K 线会无法重建当前股票池内的历史走势及市值/流动性条件。股票范围保持当前上市的沪深非 ST A 股,不包含北交所,也不维护退市股票的长期历史成员资格。价格数据只保留 Tushare `qfq` 前复权口径,不将未复权或 `hfq` 作为运行时事实。历史分析按交易日复现当前股票池内的行情与估值状态,但使用数据库中最新修订的 qfq,不承诺恢复某次同步时的 qfq 版本。
|
||||||
|
|
||||||
|
第一阶段只同步日线收盘数据,不覆盖分钟级或实时行情。
|
||||||
|
|
||||||
|
## Considered Options
|
||||||
|
|
||||||
|
- **按股票拆分 CSV 作为主存储**:保留旧项目的低门槛结构,但难以保证跨标的一致性、时点查询和增量修订。
|
||||||
|
- **SQLite 作为主存储**:适合单机原型,但与后续服务化、多进程同步和并发分析的目标不匹配。
|
||||||
|
|
||||||
|
## Consequences
|
||||||
|
|
||||||
|
- 需要设计 PostgreSQL schema、索引、迁移和备份策略。
|
||||||
|
- 同步任务必须具备批次记录、幂等 upsert 和失败重试能力。
|
||||||
|
- CSV 导出格式属于交付/快照契约,不再承担在线查询职责。
|
||||||
|
- 若未来需要未复权或 `hfq` 分析,必须重新回补数据,不能从当前主表无损推导。
|
||||||
|
- “历史可复现”不等同于恢复历史抓取版本;qfq 的供应商修订会改变历史价格,但不改变股票池和估值的交易日时点约束。
|
||||||
|
- 当前方案不覆盖长期回测所需的退市股票、历史上市状态和无幸存者偏差股票池;若需求改变,应另立数据范围决策。
|
||||||
@@ -0,0 +1,33 @@
|
|||||||
|
---
|
||||||
|
status: accepted
|
||||||
|
---
|
||||||
|
|
||||||
|
# Tushare 同步采用六年落地快照与重叠区间指纹
|
||||||
|
|
||||||
|
每日同步由外部调度器触发;应用对当前目标股票池逐只获取最近六年的 Tushare qfq 日线,形成 CSV 落地快照。新旧快照先比较完整历史重叠区间:指纹相同时只向 PostgreSQL 插入新交易日,指纹不同时对该股票最近六年执行幂等 upsert;一只股票的数据变化不会触发全市场回写。
|
||||||
|
|
||||||
|
单只股票按“生成临时 CSV → 校验并比较重叠区间 → 提交 PostgreSQL 事务 → 原子替换正式 CSV”的顺序同步。PostgreSQL 提交失败时不替换正式 CSV,保证运行时事实源和已发布落地快照不会出现新旧倒置。
|
||||||
|
|
||||||
|
单只股票也是事务和重试边界。部分股票失败时,批次标记为 `partial_success`,已成功股票不回滚;失败股票保留旧 PostgreSQL 数据和旧 CSV,并在后续只重试失败集合。
|
||||||
|
|
||||||
|
`partial_success` 可以触发后续选股,但只能使用行情和所需估值数据均已到达目标交易日的有效股票;失败股票不得使用旧数据冒充当天行情,选股结果必须携带覆盖率和失败列表。
|
||||||
|
|
||||||
|
自动选股的默认最低数据覆盖率为 `99%`,并允许通过配置调整。低于阈值时只报告同步结果并优先重试失败股票;人工强制运行必须把选股结果标记为不完整。
|
||||||
|
|
||||||
|
六年窗口按目标交易日滚动。单只股票同步事务成功时删除窗口起点之前的 PostgreSQL 日线,并以只包含当前窗口的新文件原子替换正式 CSV;每日估值/交易指标及其日期快照采用相同的滚动保留边界。失败对象不发布新 CSV,也不因本次失败提前丢失其原有数据。该边界符合当前系统关注近期走势而非长期回测的定位。
|
||||||
|
|
||||||
|
## Considered Options
|
||||||
|
|
||||||
|
- **仅请求最新增量日期**:请求次数没有明显减少,却无法自然发现历史 qfq 修订。
|
||||||
|
- **只比较最近一周**:能发现常见的近期除权变化,但可能漏掉供应商对更早历史数据的修正。
|
||||||
|
|
||||||
|
## Consequences
|
||||||
|
|
||||||
|
- CSV 是同步落地快照,PostgreSQL 仍是策略查询的事实源。
|
||||||
|
- 指纹必须基于规范化、排序后的共同日期区间和固定字段集合,避免浮点文本差异造成误判。
|
||||||
|
- PostgreSQL 写入以 Tushare `(ts_code, trade_date)` 为唯一键;未变化时只写新日期,变化时仅回写发生变化的股票。
|
||||||
|
- 批量修复应通过 staging 表与批量装载完成,避免逐行 ORM 写入。
|
||||||
|
- 临时 CSV 只有在对应 PostgreSQL 事务提交成功后才能晋升为正式落地快照;重复执行同一批次必须保持幂等。
|
||||||
|
- 同步批次需要记录每只股票的成功、失败、插入、更新和未变化结果,以支持部分成功与定向重试。
|
||||||
|
- 选股批次需要关联对应同步批次,并记录实际参与股票数、目标股票数和数据覆盖率。
|
||||||
|
- PostgreSQL 与 CSV 都只承诺滚动保留最近六年;如果未来需要更长周期回测,应重新评估数据范围并回补历史数据。
|
||||||
@@ -30,6 +30,8 @@ COPY migrations ./migrations
|
|||||||
RUN uv sync --frozen --no-dev
|
RUN uv sync --frozen --no-dev
|
||||||
|
|
||||||
RUN useradd --create-home --uid 10001 appuser
|
RUN useradd --create-home --uid 10001 appuser
|
||||||
|
RUN mkdir -p /app/data/market-data \
|
||||||
|
&& chown -R 10001:10001 /app/data/market-data
|
||||||
USER appuser
|
USER appuser
|
||||||
ENV PATH="/app/.venv/bin:$PATH"
|
ENV PATH="/app/.venv/bin:$PATH"
|
||||||
EXPOSE 8000
|
EXPOSE 8000
|
||||||
|
|||||||
@@ -2,6 +2,8 @@
|
|||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import logging
|
||||||
|
import time
|
||||||
from collections.abc import Iterable, Sequence
|
from collections.abc import Iterable, Sequence
|
||||||
from dataclasses import dataclass, field
|
from dataclasses import dataclass, field
|
||||||
from datetime import date, timedelta
|
from datetime import date, timedelta
|
||||||
@@ -15,6 +17,9 @@ from ..domain.rules import filter_current_hs_a_stocks
|
|||||||
|
|
||||||
SyncMode = Literal["daily", "initialize", "retry"]
|
SyncMode = Literal["daily", "initialize", "retry"]
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
_PROGRESS_LOG_INTERVAL = 100
|
||||||
|
|
||||||
|
|
||||||
@dataclass(frozen=True, slots=True)
|
@dataclass(frozen=True, slots=True)
|
||||||
class SyncMarketDataCommand:
|
class SyncMarketDataCommand:
|
||||||
@@ -134,8 +139,20 @@ class SyncMarketData:
|
|||||||
return self._execute_locked(command)
|
return self._execute_locked(command)
|
||||||
|
|
||||||
def _execute_locked(self, command: SyncMarketDataCommand) -> SyncBatchSummary:
|
def _execute_locked(self, command: SyncMarketDataCommand) -> SyncBatchSummary:
|
||||||
|
started_at = time.monotonic()
|
||||||
|
logger.info(
|
||||||
|
"market_data_sync_started mode=%s requested_trade_date=%s",
|
||||||
|
command.mode,
|
||||||
|
command.target_trade_date or "auto",
|
||||||
|
)
|
||||||
target_trade_date = self._resolve_target(command.target_trade_date)
|
target_trade_date = self._resolve_target(command.target_trade_date)
|
||||||
window = SyncWindow.from_target(target_trade_date)
|
window = SyncWindow.from_target(target_trade_date)
|
||||||
|
logger.info(
|
||||||
|
"market_data_sync_target target_trade_date=%s window_start=%s window_end=%s",
|
||||||
|
target_trade_date,
|
||||||
|
window.start,
|
||||||
|
window.end,
|
||||||
|
)
|
||||||
all_stocks = filter_current_hs_a_stocks(self.source.fetch_stocks())
|
all_stocks = filter_current_hs_a_stocks(self.source.fetch_stocks())
|
||||||
if not all_stocks:
|
if not all_stocks:
|
||||||
return SyncBatchSummary(
|
return SyncBatchSummary(
|
||||||
@@ -158,6 +175,11 @@ class SyncMarketData:
|
|||||||
command.parent_batch_id,
|
command.parent_batch_id,
|
||||||
len(all_stocks),
|
len(all_stocks),
|
||||||
)
|
)
|
||||||
|
logger.info(
|
||||||
|
"market_data_sync_batch_created batch_id=%s target_count=%d",
|
||||||
|
batch_id,
|
||||||
|
len(all_stocks),
|
||||||
|
)
|
||||||
failures: list[SyncFailure] = []
|
failures: list[SyncFailure] = []
|
||||||
totals = [0, 0, 0]
|
totals = [0, 0, 0]
|
||||||
stock_codes = {stock.ts_code for stock in all_stocks}
|
stock_codes = {stock.ts_code for stock in all_stocks}
|
||||||
@@ -166,18 +188,71 @@ class SyncMarketData:
|
|||||||
)
|
)
|
||||||
|
|
||||||
self._process_stock_master(batch_id, all_stocks, failures, totals)
|
self._process_stock_master(batch_id, all_stocks, failures, totals)
|
||||||
|
self._log_progress(
|
||||||
|
batch_id=batch_id,
|
||||||
|
stage="stock_master",
|
||||||
|
current=1,
|
||||||
|
total=1,
|
||||||
|
item_key="current",
|
||||||
|
totals=totals,
|
||||||
|
failures=failures,
|
||||||
|
started_at=started_at,
|
||||||
|
force=bool(failures),
|
||||||
|
)
|
||||||
dates = self._dates_to_process(window, target_trade_date, command.mode, retry_items)
|
dates = self._dates_to_process(window, target_trade_date, command.mode, retry_items)
|
||||||
for trade_date in dates:
|
logger.info(
|
||||||
if (
|
"market_data_sync_stage_started batch_id=%s stage=daily_basic total=%d",
|
||||||
command.mode == "retry"
|
batch_id,
|
||||||
and ("daily_basic", trade_date.isoformat()) not in retry_items
|
len(dates),
|
||||||
):
|
)
|
||||||
continue
|
for current, trade_date in enumerate(dates, start=1):
|
||||||
|
failure_count = len(failures)
|
||||||
self._process_daily_basic(batch_id, trade_date, stock_codes, window, failures, totals)
|
self._process_daily_basic(batch_id, trade_date, stock_codes, window, failures, totals)
|
||||||
for stock in all_stocks:
|
self._log_progress(
|
||||||
if command.mode == "retry" and ("bar", stock.ts_code) not in retry_items:
|
batch_id=batch_id,
|
||||||
continue
|
stage="daily_basic",
|
||||||
|
current=current,
|
||||||
|
total=len(dates),
|
||||||
|
item_key=trade_date.isoformat(),
|
||||||
|
totals=totals,
|
||||||
|
failures=failures,
|
||||||
|
started_at=started_at,
|
||||||
|
force=len(failures) > failure_count,
|
||||||
|
)
|
||||||
|
logger.info(
|
||||||
|
"market_data_sync_stage_completed batch_id=%s stage=daily_basic total=%d",
|
||||||
|
batch_id,
|
||||||
|
len(dates),
|
||||||
|
)
|
||||||
|
stocks_to_process = tuple(
|
||||||
|
stock
|
||||||
|
for stock in all_stocks
|
||||||
|
if command.mode != "retry" or ("bar", stock.ts_code) in retry_items
|
||||||
|
)
|
||||||
|
logger.info(
|
||||||
|
"market_data_sync_stage_started batch_id=%s stage=bar total=%d",
|
||||||
|
batch_id,
|
||||||
|
len(stocks_to_process),
|
||||||
|
)
|
||||||
|
for current, stock in enumerate(stocks_to_process, start=1):
|
||||||
|
failure_count = len(failures)
|
||||||
self._process_bar(batch_id, stock, window, failures, totals)
|
self._process_bar(batch_id, stock, window, failures, totals)
|
||||||
|
self._log_progress(
|
||||||
|
batch_id=batch_id,
|
||||||
|
stage="bar",
|
||||||
|
current=current,
|
||||||
|
total=len(stocks_to_process),
|
||||||
|
item_key=stock.ts_code,
|
||||||
|
totals=totals,
|
||||||
|
failures=failures,
|
||||||
|
started_at=started_at,
|
||||||
|
force=len(failures) > failure_count,
|
||||||
|
)
|
||||||
|
logger.info(
|
||||||
|
"market_data_sync_stage_completed batch_id=%s stage=bar total=%d",
|
||||||
|
batch_id,
|
||||||
|
len(stocks_to_process),
|
||||||
|
)
|
||||||
|
|
||||||
if not failures:
|
if not failures:
|
||||||
try:
|
try:
|
||||||
@@ -205,6 +280,20 @@ class SyncMarketData:
|
|||||||
status = "success" if not failures else "partial_success" if valid_count else "failed"
|
status = "success" if not failures else "partial_success" if valid_count else "failed"
|
||||||
eligible = coverage >= self.coverage_threshold
|
eligible = coverage >= self.coverage_threshold
|
||||||
self.repository.record_batch(batch_id, status, valid_count, coverage, eligible)
|
self.repository.record_batch(batch_id, status, valid_count, coverage, eligible)
|
||||||
|
logger.info(
|
||||||
|
"market_data_sync_finished batch_id=%s status=%s target_count=%d valid_count=%d "
|
||||||
|
"coverage=%s failures=%d inserted=%d updated=%d unchanged=%d elapsed_seconds=%.1f",
|
||||||
|
batch_id,
|
||||||
|
status,
|
||||||
|
len(all_stocks),
|
||||||
|
valid_count,
|
||||||
|
coverage,
|
||||||
|
len(failures),
|
||||||
|
totals[0],
|
||||||
|
totals[1],
|
||||||
|
totals[2],
|
||||||
|
time.monotonic() - started_at,
|
||||||
|
)
|
||||||
return SyncBatchSummary(
|
return SyncBatchSummary(
|
||||||
batch_id,
|
batch_id,
|
||||||
target_trade_date,
|
target_trade_date,
|
||||||
@@ -283,6 +372,12 @@ class SyncMarketData:
|
|||||||
self.snapshots.discard(staged)
|
self.snapshots.discard(staged)
|
||||||
failure = self._failure("stock", "current", exc)
|
failure = self._failure("stock", "current", exc)
|
||||||
failures.append(failure)
|
failures.append(failure)
|
||||||
|
logger.warning(
|
||||||
|
"market_data_sync_item_failed batch_id=%s stage=stock_master item=current "
|
||||||
|
"error_type=%s",
|
||||||
|
batch_id,
|
||||||
|
failure.error_type,
|
||||||
|
)
|
||||||
self.repository.record_item(
|
self.repository.record_item(
|
||||||
batch_id,
|
batch_id,
|
||||||
"stock",
|
"stock",
|
||||||
@@ -329,6 +424,13 @@ class SyncMarketData:
|
|||||||
self.snapshots.discard(staged)
|
self.snapshots.discard(staged)
|
||||||
failure = self._failure("daily_basic", key, exc)
|
failure = self._failure("daily_basic", key, exc)
|
||||||
failures.append(failure)
|
failures.append(failure)
|
||||||
|
logger.warning(
|
||||||
|
"market_data_sync_item_failed batch_id=%s stage=daily_basic item=%s "
|
||||||
|
"error_type=%s",
|
||||||
|
batch_id,
|
||||||
|
key,
|
||||||
|
failure.error_type,
|
||||||
|
)
|
||||||
self.repository.record_item(
|
self.repository.record_item(
|
||||||
batch_id,
|
batch_id,
|
||||||
"daily_basic",
|
"daily_basic",
|
||||||
@@ -381,6 +483,12 @@ class SyncMarketData:
|
|||||||
self.snapshots.discard(staged)
|
self.snapshots.discard(staged)
|
||||||
failure = self._failure("bar", stock.ts_code, exc)
|
failure = self._failure("bar", stock.ts_code, exc)
|
||||||
failures.append(failure)
|
failures.append(failure)
|
||||||
|
logger.warning(
|
||||||
|
"market_data_sync_item_failed batch_id=%s stage=bar item=%s error_type=%s",
|
||||||
|
batch_id,
|
||||||
|
stock.ts_code,
|
||||||
|
failure.error_type,
|
||||||
|
)
|
||||||
self.repository.record_item(
|
self.repository.record_item(
|
||||||
batch_id,
|
batch_id,
|
||||||
"bar",
|
"bar",
|
||||||
@@ -397,6 +505,41 @@ class SyncMarketData:
|
|||||||
totals[1] += result.updated
|
totals[1] += result.updated
|
||||||
totals[2] += result.unchanged
|
totals[2] += result.unchanged
|
||||||
|
|
||||||
|
@staticmethod
|
||||||
|
def _log_progress(
|
||||||
|
*,
|
||||||
|
batch_id: str,
|
||||||
|
stage: str,
|
||||||
|
current: int,
|
||||||
|
total: int,
|
||||||
|
item_key: str,
|
||||||
|
totals: list[int],
|
||||||
|
failures: Sequence[SyncFailure],
|
||||||
|
started_at: float,
|
||||||
|
force: bool = False,
|
||||||
|
) -> None:
|
||||||
|
"""Log bounded, secret-free progress for a batch stage."""
|
||||||
|
|
||||||
|
if not force and total > _PROGRESS_LOG_INTERVAL and current not in {
|
||||||
|
1,
|
||||||
|
total,
|
||||||
|
} and current % _PROGRESS_LOG_INTERVAL != 0:
|
||||||
|
return
|
||||||
|
logger.info(
|
||||||
|
"market_data_sync_progress batch_id=%s stage=%s progress=%d/%d item=%s "
|
||||||
|
"failures=%d inserted=%d updated=%d unchanged=%d elapsed_seconds=%.1f",
|
||||||
|
batch_id,
|
||||||
|
stage,
|
||||||
|
current,
|
||||||
|
total,
|
||||||
|
item_key,
|
||||||
|
len(failures),
|
||||||
|
totals[0],
|
||||||
|
totals[1],
|
||||||
|
totals[2],
|
||||||
|
time.monotonic() - started_at,
|
||||||
|
)
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
def _failure(item_kind: str, item_key: str, error: BaseException) -> SyncFailure:
|
def _failure(item_kind: str, item_key: str, error: BaseException) -> SyncFailure:
|
||||||
message = " ".join(str(error).split())[:500] or "synchronization item failed"
|
message = " ".join(str(error).split())[:500] or "synchronization item failed"
|
||||||
|
|||||||
@@ -2,6 +2,7 @@
|
|||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import logging
|
||||||
import random
|
import random
|
||||||
import time
|
import time
|
||||||
from collections.abc import Callable, Iterable, Mapping
|
from collections.abc import Callable, Iterable, Mapping
|
||||||
@@ -11,6 +12,8 @@ from typing import cast
|
|||||||
from ..domain.models import Bar, DailyBasic, Stock, SyncWindow, parse_date
|
from ..domain.models import Bar, DailyBasic, Stock, SyncWindow, parse_date
|
||||||
from ..domain.rules import filter_current_hs_a_stocks
|
from ..domain.rules import filter_current_hs_a_stocks
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
class TushareSourceError(RuntimeError):
|
class TushareSourceError(RuntimeError):
|
||||||
"""A vendor request failed after the configured retry budget."""
|
"""A vendor request failed after the configured retry budget."""
|
||||||
@@ -158,7 +161,18 @@ class TushareAdapter:
|
|||||||
if attempt == self.max_retries:
|
if attempt == self.max_retries:
|
||||||
break
|
break
|
||||||
delay = self.backoff_seconds * (2**attempt) * (0.5 + self.random_fn())
|
delay = self.backoff_seconds * (2**attempt) * (0.5 + self.random_fn())
|
||||||
|
logger.warning(
|
||||||
|
"tushare_request_retry method=%s attempt=%d max_attempts=%d",
|
||||||
|
method_name,
|
||||||
|
attempt + 1,
|
||||||
|
self.max_retries + 1,
|
||||||
|
)
|
||||||
self.sleep_fn(delay)
|
self.sleep_fn(delay)
|
||||||
|
logger.error(
|
||||||
|
"tushare_request_failed method=%s attempts=%d",
|
||||||
|
method_name,
|
||||||
|
self.max_retries + 1,
|
||||||
|
)
|
||||||
raise TushareSourceError(f"Tushare request failed: {method_name}") from last_error
|
raise TushareSourceError(f"Tushare request failed: {method_name}") from last_error
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ from __future__ import annotations
|
|||||||
|
|
||||||
import argparse
|
import argparse
|
||||||
import json
|
import json
|
||||||
|
import logging
|
||||||
from collections.abc import Sequence
|
from collections.abc import Sequence
|
||||||
from datetime import date
|
from datetime import date
|
||||||
|
|
||||||
@@ -13,6 +14,8 @@ from ..infrastructure.csv_snapshot import CsvSnapshotStore
|
|||||||
from ..infrastructure.postgres import PostgresMarketDataRepository
|
from ..infrastructure.postgres import PostgresMarketDataRepository
|
||||||
from ..infrastructure.tushare import TushareAdapter
|
from ..infrastructure.tushare import TushareAdapter
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
def build_parser() -> argparse.ArgumentParser:
|
def build_parser() -> argparse.ArgumentParser:
|
||||||
"""Build the explicit, repeatable synchronization CLI."""
|
"""Build the explicit, repeatable synchronization CLI."""
|
||||||
@@ -40,6 +43,11 @@ def main(argv: Sequence[str] | None = None) -> int:
|
|||||||
|
|
||||||
args = build_parser().parse_args(argv)
|
args = build_parser().parse_args(argv)
|
||||||
settings = get_settings()
|
settings = get_settings()
|
||||||
|
logging.basicConfig(
|
||||||
|
level=settings.log_level.upper(),
|
||||||
|
format="%(asctime)s %(levelname)s %(name)s %(message)s",
|
||||||
|
force=True,
|
||||||
|
)
|
||||||
if args.retry_batch_id:
|
if args.retry_batch_id:
|
||||||
command = SyncMarketDataCommand(
|
command = SyncMarketDataCommand(
|
||||||
mode="retry",
|
mode="retry",
|
||||||
@@ -51,6 +59,11 @@ def main(argv: Sequence[str] | None = None) -> int:
|
|||||||
mode="initialize" if args.initialize else "daily",
|
mode="initialize" if args.initialize else "daily",
|
||||||
target_trade_date=args.trade_date,
|
target_trade_date=args.trade_date,
|
||||||
)
|
)
|
||||||
|
logger.info(
|
||||||
|
"market_data_sync_cli mode=%s target_trade_date=%s",
|
||||||
|
command.mode,
|
||||||
|
command.target_trade_date or "auto",
|
||||||
|
)
|
||||||
source = TushareAdapter.from_token(
|
source = TushareAdapter.from_token(
|
||||||
settings.tushare_token,
|
settings.tushare_token,
|
||||||
max_retries=settings.market_data_max_retries,
|
max_retries=settings.market_data_max_retries,
|
||||||
|
|||||||
@@ -1,9 +1,12 @@
|
|||||||
|
import logging
|
||||||
from collections.abc import Generator, Iterable, Sequence
|
from collections.abc import Generator, Iterable, Sequence
|
||||||
from contextlib import contextmanager
|
from contextlib import contextmanager
|
||||||
from datetime import date
|
from datetime import date
|
||||||
from decimal import Decimal
|
from decimal import Decimal
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
from zhixing_server.modules.market_data.application.sync import (
|
from zhixing_server.modules.market_data.application.sync import (
|
||||||
SyncMarketData,
|
SyncMarketData,
|
||||||
SyncMarketDataCommand,
|
SyncMarketDataCommand,
|
||||||
@@ -163,3 +166,38 @@ def test_sync_is_idempotent_and_reports_coverage(tmp_path: Path) -> None:
|
|||||||
assert second.status == "success"
|
assert second.status == "success"
|
||||||
assert second.inserted_count == 1
|
assert second.inserted_count == 1
|
||||||
assert second.unchanged_count >= 2
|
assert second.unchanged_count >= 2
|
||||||
|
|
||||||
|
|
||||||
|
def test_sync_logs_progress_for_initialize_and_daily_update(
|
||||||
|
tmp_path: Path,
|
||||||
|
caplog: pytest.LogCaptureFixture,
|
||||||
|
) -> None:
|
||||||
|
target = date(2024, 1, 2)
|
||||||
|
repository = InMemoryRepository()
|
||||||
|
use_case = SyncMarketData(
|
||||||
|
FakeSource(target),
|
||||||
|
CsvSnapshotStore(tmp_path),
|
||||||
|
repository,
|
||||||
|
today=target,
|
||||||
|
)
|
||||||
|
caplog.set_level(logging.INFO, logger="zhixing_server.modules.market_data.application.sync")
|
||||||
|
|
||||||
|
initialize = use_case.execute(
|
||||||
|
SyncMarketDataCommand(mode="initialize", target_trade_date=target)
|
||||||
|
)
|
||||||
|
initialize_messages = [record.getMessage() for record in caplog.records]
|
||||||
|
|
||||||
|
assert initialize.status == "success"
|
||||||
|
assert any("stage=daily_basic" in message for message in initialize_messages)
|
||||||
|
assert any(
|
||||||
|
"stage=bar" in message and "progress=1/1" in message
|
||||||
|
for message in initialize_messages
|
||||||
|
)
|
||||||
|
|
||||||
|
caplog.clear()
|
||||||
|
daily = use_case.execute(SyncMarketDataCommand(mode="daily", target_trade_date=target))
|
||||||
|
daily_messages = [record.getMessage() for record in caplog.records]
|
||||||
|
|
||||||
|
assert daily.status == "success"
|
||||||
|
assert any("market_data_sync_started mode=daily" in message for message in daily_messages)
|
||||||
|
assert any("market_data_sync_finished" in message for message in daily_messages)
|
||||||
|
|||||||
Reference in New Issue
Block a user