融合引擎(fox_engine)
概述
融合引擎是本项目的数据写入核心,负责将多源 RAW 数据逐层归一、交叉验证、融合为 22 类 fox_* DWD 宽表,供 API 与量化层消费。整体架构遵循 RAW → STG → DWD 三层管线,以声明式注册表驱动、逐股隔离执行、增量提交防崩。
架构分层
┌─ 数据源层 ─────────────────────────────────────────────────────┐
│ 腾讯 HTTP · TDX TCP · 东财 HTTP · 新浪 HTTP · 巨潮 · 申万官方 │
└────────────────────────┬───────────────────────────────────────┘
│ 下载(逐股/批量/全市场)
▼
┌─ RAW 层 ───────────────────────────────────────────────────────┐
│ em_dragon_tiger · em_margin_trading · tdx_finance_snapshot │
│ em_block_trade · em_stock_news · sina_financial_statement ... │
│ 特点:源侧原始格式落库,不做归一,允许冗余与脏数据 │
└────────────────────────┬───────────────────────────────────────┘
│ STG Loader(14 个加载器)
▼
┌─ STG 层 ───────────────────────────────────────────────────────┐
│ stg_finance_indicator · stg_holder · stg_dividend │
│ stg_dragon_tiger · stg_margin · stg_block_trade ... │
│ 特点:类型归一 + 去重 + 质量校验,按 (code, report_date) 唯一 │
└────────────────────────┬───────────────────────────────────────┘
│ FoxFuser(22 个融合器)
▼
┌─ DWD 层 ───────────────────────────────────────────────────────┐
│ fox_stock_wide · fox_kline_daily · fox_finance_indicator │
│ fox_holder · fox_valuation · fox_industry_index_daily ... │
│ 特点:多源交叉验证,最终消费口径,API 与量化层直读 │
└────────────────────────────────────────────────────────────────┘
实现位置
| 目录 / 文件 | 职责 |
|---|---|
fox_engine/engine.py | FoxDataEngine — 读侧门面,所有 get_* 方法带分层 TTL 缓存 |
fox_engine/writer/base.py | FoxFuserSpec 声明 + run_fuser() 通用执行器 |
fox_engine/writer/registry.py | FOX_FUSERS 注册表(22 个 FoxFuserSpec) |
fox_engine/writer/freshness.py | 股票级新鲜度过滤(跳过已达最新交易日的股票) |
fox_engine/writer/fusers/ | 各数据类型的融合实现(kline_daily_fox.py / stock_wide.py 等) |
fox_engine/writer/fusers/_raw_heal.py | RAW 过期自愈(融合前自动补拉) |
fox_engine/writer/reconcile.py | 后置对账/清账(成交额派生、快照残留清理、周期K去重) |
fox_engine/stg_loader/ | 14 个 STG 加载器(RAW → STG 归一) |
fox_engine/stg_loader/base.py | StgLoaderSpec 声明 + run_stg_loader() 执行器 |
fox_engine/cache.py | BoundedTTLCache — 进程级单飞 LRU 缓存 |
声明式融合注册表
每个 fox_* 表对应一个 FoxFuserSpec,声明式注册到 FOX_FUSERS 字典:
FoxFuserSpec(
data_type="finance_indicator", # 唯一标识
target_table=FoxFinanceIndicator, # ORM 模型
read_sources=["sina_financial_indicator", "tdx_finance_snapshot"],
list_codes=list_all_codes, # 股票池枚举函数
build_rows=build_finance_indicator, # 单股行构建函数
freshness=FreshnessSpec(column="updated_at", stale_days=1),
raw_heal=RawHealSpec(raw_table="sina_financial_statement", stale_days=2),
incremental_commit=True, # 大表增量提交
fill_only_cols=("eps", "bvps"), # 缺失补位列(0/NULL 不覆盖非 0 旧值)
)
新增 fox_ 表 = 一个 ORM 模型 + 一条注册表项,无需改执行器。
融合执行循环(run_fuser)
1. RAW 自愈预检(若声明了 RawHealSpec + 全市场口径)
2. STG 预加载(RAW → STG via stg_loader)
3. 开启任务(FoxSyncTask 行,带 worker_id)
4. STG 预检(源表为空 → quality warning)
5. 读源 → 枚举代码 → 新鲜度过滤(跳过已新鲜股票)
6. 逐股构建循环(每 200 只回报进度)
├─ 熔断:失败率 > 30%(≥3 样本后)→ 整批中止
├─ 增量提交:大表每 commit_every_codes 只 flush 一次
└─ 空交付守卫:全市场 0 行 → failed(防止 last_success 被误标)
7. 质量摘要(空值率 + 过期行检查 → fox_data_quality)
8. 血缘记录(FoxDataRecord 带 batch_id / latency_ms)
9. 僵尸清扫(sweep_zombie_tasks 检测死亡进程)
关键设计决策
| 决策 | 原因 |
|---|---|
| 逐股隔离 | 单股异常不污染整批;失败率熔断保护下游 |
| 增量提交 | fox_kline_daily 5500+ 只全量内存占用过大,每 500 只 flush |
| 空交付判 failed | fox_stock_wide 曾因"报成功但 0 行"静默停更 16 天 |
| 新鲜度只在全市场生效 | len(codes) >= 1000 才启用,单股/小批量始终重拉 |
| fill_only_cols | 融合器取不到的列用 COALESCE(NULLIF(新值,0), 旧值) 缺失补位,不覆盖已有非 0 值 |
新鲜度过滤(freshness.py)
filter_fresh_codes() 查询目标表,跳过已有最新交易日数据的股票:
- 阈值:
latest_trade_date()00:00(15:30 翻日,不是current_or_last_trading_day()的 09:30) - 盘后安全阀:当数据可用日 == 今天时,额外要求
updated_at >= 15:00(防止盘中快照被冻结为收盘价) - 范围守卫:仅
len(target_codes) > 1000时启用(子集/单股始终重拉)
RAW 自愈(_raw_heal.py)
ensure_raw_fresh() 在融合前检查 RAW 表 MAX(date_col):
- 触发条件:RAW 过期 ≥
stale_days(默认 2 天)+ MySQL + 全市场口径 - 自愈动作:对自选股逐股下载(随机节流)→ 重跑对应 STG 加载器
- 声明位置:6 个事件类融合器(holder / dividend / dragon_tiger / margin / block_trade / announcement)
- 不触发自愈:API 直拉型融合器(kline / valuation / stock_wide)直接拉 API,不走 RAW
STG 加载器(stg_loader/)
14 个加载器将 RAW 归一为 STG,与 DWD 融合器共用同一模式(read_raw / list_codes / build_rows):
| STG 表 | 来源 RAW | 加载器 |
|---|---|---|
stg_finance_indicator | sina_financial_statement + tdx_finance_snapshot | finance_indicator.py |
stg_holder | em_holder + tdx_finance_snapshot | holder.py |
stg_dividend | em_dividend | dividend.py |
stg_dragon_tiger | em_dragon_tiger + ths_f10_lhb | dragon_tiger.py |
stg_margin | em_margin_trading | margin.py |
stg_block_trade | em_block_trade + ths_f10_dzjy | block_trade.py |
stg_announcement | cninfo_announcement | announcement.py |
stg_lockup | em_lockup_expiry | lockup.py |
stg_stock_master | tdx_finance_snapshot | stock_master.py |
| ... | ... | ... |
5 类跳过 STG(API 直拉):kline_daily_fox / kline_bar / index_daily / stock_wide / market_dpyt
事件表 truncate-reload:事件 STG 表(自增 PK、无业务唯一键)采用 DELETE-all + re-INSERT 防止重复累积。
后置对账(reconcile.py)
三个幂等、失败安全的后置步骤,不阻塞主融合:
| 函数 | 作用 | 调用方 |
|---|---|---|
reconcile_amount_derived | 用 tencent_stock_quote_daily 当日成交额补齐 fox_kline_daily.amount 与 fox_stock_wide/fox_fund_flow_daily 的 main_net_ratio_pct(按日 JOIN) | fox_wide_sync |
prune_snapshot_residue | 清理 fox_stock_wide 残留行(段外 B 股 + 停更 A 股:最后 K 线 > 120 天且当日快照 ≥ 3000 只) | fox_wide_sync |
purge_period_residue | 周期 K 线去重(同一 (code, period, 周期桶) 只留 MAX(date)) | fox_kline_sync |
reconcile_research_coverage | 从 em_research_report 90 日窗口回填 fox_stock_wide 研报列(目标价/评级) | run_reconcile |
三域任务拓扑
从单体 sync_all 拆分为 3 个独立域任务(类型不相交、各自独立运行锁、可并行):
| 域任务 | 时间 | 融合类型 | 后置动作 |
|---|---|---|---|
fox_kline_sync | 15:40 | kline_daily_fox + kline_bar(2 类) | purge_period_residue |
fox_daily_sync | 16:00 | market_daily / market_dpyt / valuation / index_daily / company_profile / finance_indicator / holder / dividend / lockup / dragon_tiger / margin / block_trade / announcement / news(14 类) | — |
fox_wide_sync | 16:20 | stock_master / industry / stock_wide / industry_index(4 类) | run_reconcile(amount_derived + prune_snapshot + research_coverage) |
依赖关系:
fox_kline_sync→quant_calc(17:30 读fox_kline_daily)fox_daily_sync事件类 →evening_sync先写 RAW(或 RAW 自愈兜底)fox_wide_sync→selection_sync(18:00 读fox_stock_wide)
读侧降级链
以 K 线为例(stock_service.get_kline),四层降级:
1. SQLite 本地快路径(零网络,新鲜度 + 窗口守卫)
↓ miss
2. KlineStore Parquet 镜像(本地列存,零 MySQL 开销)
↓ miss
3. MySQL fox_kline_daily(DWD 权威源,新鲜度检查)
↓ miss
4. 实时拉取(TDX → 腾讯 → Baostock → 东财)→ 回写 SQLite + MySQL
"优先读当日落库"原则:融合器取行情类数据时第一读必须是当日已落库的 DWD/时序表(如 tencent_stock_quote_daily),命中即返回;落库缺失时就地自愈落库再读;只有落库渠道彻底不可用才直连 API 并记 warning。
fox_ 表完整清单(22 类)
| fox_ 表 | 融合类型 | 数据源 | 域任务 |
|---|---|---|---|
fox_stock_master | stock_master | TDX + 巨潮 | fox_wide_sync |
fox_company_profile | company_profile | F10 概览 | fox_daily_sync |
fox_finance_indicator | finance_indicator | 新浪 + TDX | fox_daily_sync |
fox_valuation | valuation | 腾讯 quote | fox_daily_sync |
fox_industry | industry | TDX + STG | fox_wide_sync |
fox_holder | holder | 东财 + TDX | fox_daily_sync |
fox_dividend | dividend | 东财 | fox_daily_sync |
fox_lockup_expiry | lockup | 东财 | fox_daily_sync |
fox_dragon_tiger | dragon_tiger | 东财 + THS | fox_daily_sync |
fox_margin | margin | 东财 | fox_daily_sync |
fox_block_trade | block_trade | 东财 | fox_daily_sync |
fox_announcement | announcement | 巨潮 | fox_daily_sync |
fox_news | news | 东财 + THS | fox_daily_sync |
fox_kline_daily | kline_daily_fox | 腾讯 HTTP | fox_kline_sync |
fox_index_daily | index_daily | 腾讯 HTTP | fox_daily_sync |
fox_kline_bar | kline_bar | 腾讯 HTTP | fox_kline_sync |
fox_stock_wide | stock_wide | TDX + 腾讯 | fox_wide_sync |
fox_market_daily | market_daily | JRJ + TDX | fox_daily_sync |
fox_market_dpyt | market_dpyt | JRJ dpyt | fox_daily_sync |
fox_industry_index_daily | industry_index | 申万官方 | fox_wide_sync |
fox_etf_daily | etf_daily | 交易所 + 腾讯 | etf_daily_sync |
fox_fund_flow_daily | fund_flow | TDX + 东财兜底 | fund_flow_sync |
ODS 辅助层(计算产物,非融合表)
定义于 backend/app/models/fox_models_ods.py(15 张表),与 DWD 融合表共用 fox_ 前缀,但语义完全不同:
| ODS 辅助表 | 产出任务 | 输入源 | 说明 |
|---|---|---|---|
fox_stock_indicators_daily | quant_calc | fox_kline_daily | 技术指标(RSI/KDJ/MACD/BOLL 等) |
fox_stock_pattern_daily | quant_calc | fox_kline_daily | K线形态 + 结构形态 |
fox_stock_cyq_daily | quant_calc | fox_kline_daily | 筹码分布 |
fox_valuation_percentile | quant_calc | tencent_stock_quote_daily | PE-TTM/PB 历史分位 |
fox_stock_fscore | fscore_calc | sina_financial_statement | Piotroski F-Score |
fox_float_share_daily | quant_calc | TDX 财务快照 | 流通股本变动 |
fox_industry_prosperity_daily | quant_calc | fox_stock_wide | 行业景气度 |
fox_quote_snapshot | 实时采集 | 腾讯/TDX | 盘中实时行情快照 |
fox_fund_flow_daily | fund_flow_sync | TDX + 东财 | 日频资金流向 |
fox_strategy_pick_daily | strategy_auto_run | fox_stock_wide + indicators | 策略选股结果 |
fox_indicator_signal_pool | quant_calc | fox_stock_indicators_daily | 技术指标信号池 |
fox_stock_f10_info | fox_daily_sync | F10 概览 | 公司概况缓存 |
fox_stock_left_daily | quant_calc | fox_kline_daily | 左Bar数据 |
fox_stock_mispricing_daily | quant_calc | fox_stock_wide | 定价偏差 |
fox_stock_trade_plan_daily | 策略任务 | 多源 | 每日交易计划 |
为何绕过 STG 直写
这些表的输入源单一且确定(绝大多数只读 fox_kline_daily 或 fox_stock_wide),计算逻辑是确定性规则(RSI 公式、形态判据、F-Score 评分),不存在多源冲突需要归一化或交叉验证——走 STG 层是过度设计。
与 NO_ODS_TYPES(8 类 DWD 融合器绕过 STG 直拉 API)的区别:
| 维度 | NO_ODS_TYPES | ODS 辅助层 |
|---|---|---|
| 目标表 | fox_* DWD 融合表 | fox_* 计算产物表(fox_models_ods.py) |
| 绕过方向 | 跳过 STG,融合器直拉 API | 跳过 STG,计算任务直写 |
| 原因 | API 已返回归一化数据,无 RAW 可落 | 输入源单一、计算确定性高,无多源冲突 |
| 注册位置 | registry.py::NO_ODS_TYPES | fox_models_ods.py 模型定义 |
前缀共用说明
ODS 辅助层与 DWD 融合层共用 fox_ 前缀,区分靠注册表(FOX_FUSERS 只登记 22 类融合表,ODS 表不在其中)和模型文件(fox_models.py vs fox_models_ods.py)。这是历史演进的产物——早期所有 fox_ 表都归融合引擎管理,后来量化层产出表增多,单独拆出 fox_models_ods.py 但保留了前缀。
如需新增量化/计算类产出表,放入
fox_models_ods.py,不要加入FOX_FUSERS注册表。
注意事项
- 三域独立锁:
fox_kline_sync/fox_daily_sync/fox_wide_sync各自独立运行锁,可并行;重跑/接力时股票级新鲜度过滤跳过已达最新交易日数据的股票(增量续传) - 空交付 ≠ 成功:候选非空却 0 行记
failed+ 质量记录,防止last_success_date被误标 - 盘中不写落库表:
load_tencent_quotes等入口一律加day <= latest_trade_date()门禁,未收盘只走直连、不读不写落库表(防止盘中价被冻结为收盘价) - fill_only_cols 语义:融合器取不到的列用
COALESCE(NULLIF(新值,0), 旧值),0/NULL 视同未取到,非 0 不覆盖——防止整段重拉清空别源补的列 - RAW 自愈不替代定时任务:自愈只对自选股逐股补拉,全市场覆盖仍靠
wholemarket_events_sync18:50 批量刷新 - STG truncate-reload:事件 STG 表无业务唯一键,DELETE-all + re-INSERT 是唯一安全的去重方式
相关文档
- 数据源总览
- 数据流转
- 数据模型
- 数据契约(分层/单位/血缘)
docs/reports/data-coverage-backfill-2026-10-07.md(数据覆盖与回补报告) — 22 类 fox_ 融合表最新覆盖日期与缺口