数据处理管道

📎 引用文件

本文引用的文件 - agent/backtest/models.py - agent/src/entities/models.py - agent/src/entities/cashflow.py - agent/src/entities/ingest.py - agent/backtest/loaders/base.py - agent/src/factors/base.py - agent/src/factors/_backend.py - agent/src/factors/registry.py

目录

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

简介

本文件为 Vibe-Trading 数据处理管道的权威文档,覆盖数据摄入、转换规则与验证机制;详述核心数据模型(Position、FillRecord、TradeRecord、EquitySnapshot)的设计与使用;解释因子计算管道、技术指标处理与数据聚合方法;说明数据缓存策略、内存管理与性能优化技术;提供数据质量检查、异常处理与回滚机制;并以端到端工作流串联从原始数据到分析结果的完整转换过程。

项目结构

Vibe-Trading 的数据处理由两条并行路径组成: - 非条形(non-bar)现金流与面板数据路径:用于基金、债券、事件流等非 OHLCV 的时序数据,强调严格的货币、日期与符号约定。 - 条形(bar)数据加载与回测路径:面向 OHLCV 行情数据,提供统一的数据加载协议、校验、重试与本地缓存。

graph TB subgraph "非条形数据" A["entities/ingest.py<br/>load_cashflows / load_panel"] --> B["entities/cashflow.py<br/>CashFlow / CashFlowSeries / FX"] C["entities/models.py<br/>Entity / Instrument / Security / Fund / Bond"] end subgraph "条形数据" D["backtest/loaders/base.py<br/>校验/重试/缓存"] --> E["各市场 Loader<br/>yahoo/tushare/ccxt/..."] end subgraph "因子与指标" F["factors/base.py<br/>rank/zscore/ts_* 算子"] --> G["factors/registry.py<br/>Alpha 注册与执行"] end subgraph "回测模型" H["backtest/models.py<br/>Position/FillRecord/TradeRecord/EquitySnapshot"] end A --> B C --> B D --> E E --> G G --> H

图表来源 - agent/src/entities/ingest.py:316-486 - agent/src/entities/cashflow.py:121-423 - agent/src/entities/models.py:141-397 - agent/backtest/loaders/base.py:31-119 - agent/src/factors/base.py:62-356 - agent/src/factors/registry.py:201-397 - agent/backtest/models.py:13-118

章节来源 - agent/src/entities/ingest.py:1-20 - agent/backtest/loaders/base.py:1-8

核心组件

章节来源 - agent/src/entities/models.py:68-138 - agent/src/entities/cashflow.py:62-87 - agent/src/entities/ingest.py:72-135 - agent/backtest/loaders/base.py:31-119 - agent/backtest/loaders/base.py:163-236 - agent/backtest/loaders/base.py:243-439 - agent/src/factors/base.py:62-356 - agent/src/factors/_backend.py:1-69 - agent/src/factors/registry.py:87-125 - agent/backtest/models.py:13-118

架构总览

下图展示从原始数据到因子与回测结果的端到端流水线。

sequenceDiagram participant U as "用户/上游" participant ING as "ingest.py<br/>load_cashflows/load_panel" participant CF as "cashflow.py<br/>CashFlow/CashFlowSeries" participant LDR as "loaders/base.py<br/>校验/重试/缓存" participant FAC as "factors/base.py<br/>算子" participant REG as "factors/registry.py<br/>Registry.compute" participant BT as "backtest/models.py<br/>Position/Fill/Trade/Equity" U->>ING : 读取现金流/面板文件 ING-->>CF : 构建不可变序列/面板 U->>LDR : 拉取OHLCV(多源) LDR-->>U : 返回DataFrame(带校验/缓存命中) U->>FAC : 计算技术指标/横截面标准化 FAC-->>REG : 组装panel字典 REG-->>U : 输出因子矩阵(NaN传播/无inf) U->>BT : 生成头寸/成交/权益快照

图表来源 - agent/src/entities/ingest.py:316-486 - agent/src/entities/cashflow.py:196-423 - agent/backtest/loaders/base.py:31-119 - agent/backtest/loaders/base.py:343-439 - agent/src/factors/base.py:62-356 - agent/src/factors/registry.py:321-397 - agent/backtest/models.py:13-118

详细组件分析

数据摄入:现金流与面板

flowchart TD S["开始"] --> R["读取CSV/TSV/TXT"] R --> M["列名映射(默认别名/显式覆盖)"] M --> V1{"date/amount 存在?"} V1 -- 否 --> E1["抛出入参错误(指明缺失列)"] V1 -- 是 --> P["解析金额(符号/分组/小数点)"] P --> D["解析日期(ISO或format)"] D --> K{"kind/currency 存在?"} K -- 否 --> E2["抛出入参错误(指明缺失字段)"] K -- 是 --> O["构建不可变对象(CashFlow/PanelObservation)"] O --> T["排序/聚合(CashFlowSeries/EntityPanel)"] T --> End["结束"]

图表来源 - agent/src/entities/ingest.py:76-135 - agent/src/entities/ingest.py:138-265 - agent/src/entities/ingest.py:316-486 - agent/src/entities/ingest.py:538-691

章节来源 - agent/src/entities/ingest.py:49-135 - agent/src/entities/ingest.py:138-265 - agent/src/entities/ingest.py:316-486 - agent/src/entities/ingest.py:538-691

非条形数据模型:现金流与汇率

classDiagram class CashFlow { +date +amount +kind +currency +metadata +is_valuation() bool } class CashFlowSeries { +flows +currency +pre_translated +filter(kind) +between(start,end) +total(include_valuations) } class FxRate { +base_currency +quote_currency +date +rate } class FxRateTable { +quote_currency +rates +get_rate(base,on_date,allow_stale,max_staleness_days) } CashFlowSeries --> CashFlow : "包含" FxRateTable --> FxRate : "索引"

图表来源 - agent/src/entities/cashflow.py:121-194 - agent/src/entities/cashflow.py:196-423 - agent/src/entities/cashflow.py:443-681 - agent/src/entities/cashflow.py:684-769

章节来源 - agent/src/entities/cashflow.py:62-87 - agent/src/entities/cashflow.py:121-194 - agent/src/entities/cashflow.py:196-423 - agent/src/entities/cashflow.py:443-681 - agent/src/entities/cashflow.py:684-769

条形数据加载与校验

flowchart TD Q["请求fetch(source,symbol,...)"] --> CK{"缓存启用且区间已结算?"} CK -- 否 --> NET["调用外部API(带重试/预算)"] CK -- 是 --> READ["读取parquet+元数据"] READ --> HIT{"命中且有效?"} HIT -- 是 --> RET["返回缓存DataFrame"] HIT -- 否 --> NET NET --> RES{"返回DataFrame?"} RES -- 否 --> ERR["抛出TimeoutError/原始异常"] RES -- 是 --> PUT{"可缓存? (非空/已结算)"} PUT -- 是 --> WRITE["原子写入parquet+json"] PUT -- 否 --> RET2["返回DataFrame"] WRITE --> RET2

图表来源 - agent/backtest/loaders/base.py:31-119 - agent/backtest/loaders/base.py:163-236 - agent/backtest/loaders/base.py:243-439 - agent/backtest/loaders/base.py:475-596

章节来源 - agent/backtest/loaders/base.py:31-119 - agent/backtest/loaders/base.py:163-236 - agent/backtest/loaders/base.py:243-439 - agent/backtest/loaders/base.py:475-596

因子计算管道与指标处理

sequenceDiagram participant U as "调用方" participant REG as "Registry.compute" participant MOD as "Alpha模块" participant OPS as "base算子" U->>REG : 传入 alpha_id 与 panel REG->>REG : 校验required/extras/sector REG->>MOD : 惰性import并获取compute MOD->>OPS : 调用rank/zscore/ts_*等 OPS-->>MOD : 返回因子矩阵(NaN传播/无inf) MOD-->>REG : 返回DataFrame REG->>REG : 校验形状/inf/nan_ratio REG-->>U : 返回因子结果

图表来源 - agent/src/factors/base.py:62-356 - agent/src/factors/_backend.py:1-69 - agent/src/factors/registry.py:87-125 - agent/src/factors/registry.py:321-397

章节来源 - agent/src/factors/base.py:1-13 - agent/src/factors/base.py:62-356 - agent/src/factors/_backend.py:1-69 - agent/src/factors/registry.py:201-397

回测核心数据模型

classDiagram class Position { +symbol +direction +entry_price +entry_time +size +leverage +entry_bar_idx +entry_commission } class FillRecord { +symbol +timestamp +bar_idx +action +signed_quantity +notional +execution_price +fee +margin +reason +holding_bars } class TradeRecord { +symbol +direction +entry_price +exit_price +entry_time +exit_time +size +leverage +pnl +pnl_pct +exit_reason +holding_bars +commission +entry_margin +exit_margin } class EquitySnapshot { +timestamp +capital +unrealized +equity +positions }

图表来源 - agent/backtest/models.py:13-118

章节来源 - agent/backtest/models.py:13-118

依赖关系分析

graph LR INGEST["ingest.py"] --> CF["cashflow.py"] INGEST --> MODELS["models.py"] LOADERS["loaders/base.py"] --> |校验/重试/缓存| FACTORS["factors/base.py"] FACTORS --> REG["factors/registry.py"] REG --> MODELS_BT["backtest/models.py"]

图表来源 - agent/src/entities/ingest.py:316-486 - agent/src/entities/cashflow.py:196-423 - agent/backtest/loaders/base.py:343-439 - agent/src/factors/base.py:62-356 - agent/src/factors/registry.py:321-397 - agent/backtest/models.py:13-118

章节来源 - agent/src/entities/ingest.py:316-486 - agent/backtest/loaders/base.py:343-439 - agent/src/factors/base.py:62-356 - agent/src/factors/registry.py:321-397 - agent/backtest/models.py:13-118

性能考虑

[本节为通用指导,无需特定文件引用]

故障排查指南

章节来源 - agent/src/entities/ingest.py:138-265 - agent/src/entities/cashflow.py:443-681 - agent/backtest/loaders/base.py:31-119 - agent/backtest/loaders/base.py:163-236 - agent/backtest/loaders/base.py:343-439 - agent/src/factors/base.py:306-356 - agent/src/factors/registry.py:376-397 - agent/backtest/models.py:38-118

结论

Vibe-Trading 的数据处理管道以“强约束、可追溯、高性能”为核心设计原则: - 非条形数据路径通过严格的符号、币种与日期规范,确保现金流与面板数据的正确性与可审计性。 - 条形数据路径提供统一的加载协议、健壮的质量校验、可控的重试与高效的本地缓存。 - 因子计算管道以可组合的算子与注册表机制实现灵活扩展,并通过后端加速保障性能。 - 回测模型以不可变数据类固化交易证据与权益快照,支撑可靠的分析与复盘。

[本节为总结性内容,无需特定文件引用]

附录

[本节为概念性说明,无需特定文件引用]