数据处理管道

📎 引用文件

本文引用的文件 - base.py - cn_adjust.py - _symbol_utils.py - yahoo_loader.py - tushare.py - akshare_loader.py - ingest.py - correlation.py

目录

  1. 简介
  2. 项目结构
  3. 核心组件
  4. 架构总览
  5. 详细组件分析
  6. 依赖关系分析
  7. 性能考虑
  8. 故障排查指南
  9. 结论
  10. 附录:使用示例

简介

本文件面向 Vibe-Trading 的数据处理管道,聚焦于“数据摄入—清洗转换—标准化输出”的端到端流程。文档重点覆盖: - 原始数据获取:多来源(A 股、美股、港股、加密货币、外汇等)统一接入与路由 - 清洗与标准化:OHLC 校验、时间索引对齐、字段归一化、复权因子应用(中国市场特殊处理) - 符号解析与映射:跨市场股票代码的统一规范(如 AAPL.US、0700.HK、600000.SH) - 批量优化:并行/重试/预算控制、本地缓存、内存友好写入 - 质量监控:统计与异常检测思路、可观测性建议 - 实际使用示例:从不同来源拉取并转换为统一格式

项目结构

数据处理管道主要由以下模块构成: - 加载器协议与通用工具:定义 DataLoader 协议、日期/OHLC 校验、重试与预算、本地缓存 - 市场特定加载器:Yahoo Finance、Tushare、AKShare 等 - 中国 A 股复权:前复权因子计算与应用 - 非 OHLC 实体面板:现金流与面板数据的稳健导入 - 符号解析:统一符号到各加载器期望形式

graph TB subgraph "加载层" LBase["加载器基类与工具<br/>base.py"] LYahoo["Yahoo 加载器<br/>yahoo_loader.py"] LTushare["Tushare 加载器<br/>tushare.py"] LAKShare["AKShare 加载器<br/>akshare_loader.py"] end subgraph "标准化层" SAdj["中国复权因子<br/>cn_adjust.py"] SSym["符号工具<br/>_symbol_utils.py"] SCorr["符号解析/映射<br/>correlation.py"] end subgraph "非OHLC数据" EIngest["现金流/面板导入<br/>ingest.py"] end LBase --> LYahoo LBase --> LTushare LBase --> LAKShare LTushare --> SAdj LAKShare --> SAdj LYahoo --> SCorr LTushare --> SCorr LAKShare --> SCorr EIngest -.-> LBase

图表来源 - base.py:1-645 - cn_adjust.py:1-78 - _symbol_utils.py:1-21 - yahoo_loader.py:1-200 - tushare.py:1-200 - akshare_loader.py:1-200 - correlation.py:57-93 - ingest.py:1-800

章节来源 - base.py:1-645 - cn_adjust.py:1-78 - _symbol_utils.py:1-21 - yahoo_loader.py:1-200 - tushare.py:1-200 - akshare_loader.py:1-200 - correlation.py:57-93 - ingest.py:1-800

核心组件

章节来源 - base.py:27-120 - base.py:163-236 - base.py:243-439 - cn_adjust.py:27-78 - yahoo_loader.py:31-106 - tushare.py:18-80 - akshare_loader.py:20-72 - _symbol_utils.py:12-21 - correlation.py:57-93

架构总览

下图展示从“用户请求”到“标准化 DataFrame”的调用链,以及关键的质量控制点(日期校验、OHLC 校验、复权、缓存)。

sequenceDiagram participant U as "调用方" participant R as "注册表/路由" participant L as "具体加载器" participant B as "基础工具(base)" participant C as "缓存(cached_loader_fetch)" participant A as "复权(cn_adjust)" U->>R : fetch(codes, start, end, interval) R->>L : 选择匹配市场的 DataLoaders L->>B : validate_date_range() L->>C : cached_loader_fetch(...) alt 命中缓存 C-->>L : DataFrame(已标准化) else 未命中 L->>L : 调用外部API获取原始数据 L->>B : validate_ohlc() opt 中国市场(Tushare/AKShare) L->>A : apply_qfq(df, factor) A-->>L : 前复权后的DataFrame end L->>C : 写入缓存(仅最终区间) C-->>L : 返回DataFrame end L-->>U : {symbol : DataFrame}

图表来源 - base.py:163-236 - base.py:243-439 - cn_adjust.py:27-78 - tushare.py:139-200 - akshare_loader.py:93-156 - yahoo_loader.py:125-170

详细组件分析

组件A:数据加载器协议与通用工具

flowchart TD Start(["进入加载"]) --> VDate["校验日期范围"] VDate --> CacheHit{"缓存命中?"} CacheHit -- 是 --> ReturnCache["返回缓存DataFrame"] CacheHit -- 否 --> Fetch["调用外部API"] Fetch --> ValidateOHLC["OHLC不变量校验"] ValidateOHLC --> ApplyAdj{"需要复权?"} ApplyAdj -- 是 --> QFQ["apply_qfq 前复权"] ApplyAdj -- 否 --> Normalize["时间索引/字段标准化"] QFQ --> Normalize Normalize --> PutCache["写入缓存(仅最终区间)"] PutCache --> End(["返回DataFrame"])

图表来源 - base.py:31-120 - base.py:243-439 - cn_adjust.py:27-78

章节来源 - base.py:27-120 - base.py:163-236 - base.py:243-439

组件B:中国市场复权因子(前复权)

flowchart TD S(["输入: 原始日线 + adj_factor"]) --> Check["因子存在且有效?"] Check -- 否 --> RetNone["返回 None"] Check -- 是 --> Reindex["重采样到交易日"] Reindex --> Ratio["ratio = factor / factor[-1]"] Ratio --> PriceAdj["价格列 *= ratio"] Ratio --> VolAdj["volume /= ratio"] PriceAdj --> Out["输出: 前复权DataFrame"] VolAdj --> Out

图表来源 - cn_adjust.py:27-78

章节来源 - cn_adjust.py:27-78

组件C:符号解析与映射(多市场统一)

flowchart TD In["输入: code, market"] --> Crypto{"crypto 或已有 '.' ?"} Crypto -- 是 --> Pass["原样返回"] Crypto -- 否 --> Upper["大写化"] Upper --> US{"us_equity?"} US -- 是 --> AddUS["追加 .US"] US -- 否 --> HK{"hk_equity?"} HK -- 是 --> AddHK["追加 .HK"] HK -- 否 --> CN{"a_share?"} CN -- 是 --> Judge["按首数字判断 SH/BJ/SZ"] CN -- 否 --> Keep["保持原码"]

图表来源 - correlation.py:57-93

章节来源 - correlation.py:57-93

组件D:Yahoo 加载器(全球权益)

章节来源 - yahoo_loader.py:1-200

组件E:Tushare 加载器(A 股/港股/期货/基金)

章节来源 - tushare.py:1-200

组件F:AKShare 加载器(A 股/美股/港股/外汇/宏观)

章节来源 - akshare_loader.py:1-200

组件G:非 OHLC 实体面板(现金流/指标面板)

章节来源 - ingest.py:1-800

依赖关系分析

graph LR Base["base.py"] --> Yahoo["yahoo_loader.py"] Base --> Tushare["tushare.py"] Base --> AKShare["akshare_loader.py"] Tushare --> Adj["cn_adjust.py"] AKShare --> Adj Corr["correlation.py"] --> Yahoo Corr --> Tushare Corr --> AKShare

图表来源 - base.py:1-645 - cn_adjust.py:1-78 - yahoo_loader.py:1-200 - tushare.py:1-200 - akshare_loader.py:1-200 - correlation.py:57-93

章节来源 - base.py:1-645 - cn_adjust.py:1-78 - yahoo_loader.py:1-200 - tushare.py:1-200 - akshare_loader.py:1-200 - correlation.py:57-93

性能考虑

[本节为通用指导,不直接分析具体文件]

故障排查指南

章节来源 - base.py:31-120 - base.py:163-236 - base.py:243-439 - tushare.py:18-80 - cn_adjust.py:27-78 - ingest.py:316-486

结论

Vibe-Trading 的数据处理管道通过统一的加载器协议、严格的校验与复权机制、以及健壮的本地缓存与重试策略,实现了跨市场、多来源数据的稳定摄入与标准化。中国市场的前复权处理消除了除权日的机械跳空,保证了收益序列的正确性。对于大规模批量处理,建议充分利用缓存与预算控制,并结合日志与质量监控持续优化。

[本节为总结,不直接分析具体文件]

附录:使用示例

以下为典型工作流说明(不展示代码片段,仅提供路径参考): - 示例1:拉取 A 股日线并前复权 - 步骤:选择 tushare 加载器 → 校验日期 → 拉取日线 → 合并基本面(可选)→ 应用前复权 → 写入缓存 - 参考路径:tushare.py:139-200、cn_adjust.py:27-78 - 示例2:拉取美股/港股日线 - 步骤:选择 yahoo 加载器 → 区间映射 → 构建 DataFrame → 标准化时间索引 → 校验 OHLC - 参考路径:yahoo_loader.py:125-170 - 示例3:拉取 A 股/美股/港股/外汇(AKShare) - 步骤:识别市场类型(ETF/A 股/美股/港股/外汇)→ 调用对应接口 → 标准化列与日期 - 参考路径:akshare_loader.py:93-156 - 示例4:导入现金流/面板文件 - 步骤:指定列映射/货币/单位/日期格式 → 解析金额与日期 → 组装 CashFlowSeries/EntityPanel - 参考路径:ingest.py:316-486、ingest.py:694-800

[本节为使用指引,不直接分析具体文件]