跳到主要内容

融合引擎(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.pyFoxDataEngine — 读侧门面,所有 get_* 方法带分层 TTL 缓存
fox_engine/writer/base.pyFoxFuserSpec 声明 + run_fuser() 通用执行器
fox_engine/writer/registry.pyFOX_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.pyRAW 过期自愈(融合前自动补拉)
fox_engine/writer/reconcile.py后置对账/清账(成交额派生、快照残留清理、周期K去重)
fox_engine/stg_loader/14 个 STG 加载器(RAW → STG 归一)
fox_engine/stg_loader/base.pyStgLoaderSpec 声明 + run_stg_loader() 执行器
fox_engine/cache.pyBoundedTTLCache — 进程级单飞 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
空交付判 failedfox_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_indicatorsina_financial_statement + tdx_finance_snapshotfinance_indicator.py
stg_holderem_holder + tdx_finance_snapshotholder.py
stg_dividendem_dividenddividend.py
stg_dragon_tigerem_dragon_tiger + ths_f10_lhbdragon_tiger.py
stg_marginem_margin_tradingmargin.py
stg_block_tradeem_block_trade + ths_f10_dzjyblock_trade.py
stg_announcementcninfo_announcementannouncement.py
stg_lockupem_lockup_expirylockup.py
stg_stock_mastertdx_finance_snapshotstock_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_sync15:40kline_daily_fox + kline_bar(2 类)purge_period_residue
fox_daily_sync16:00market_daily / market_dpyt / valuation / index_daily / company_profile / finance_indicator / holder / dividend / lockup / dragon_tiger / margin / block_trade / announcement / news(14 类)—
fox_wide_sync16:20stock_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_masterstock_masterTDX + 巨潮fox_wide_sync
fox_company_profilecompany_profileF10 概览fox_daily_sync
fox_finance_indicatorfinance_indicator新浪 + TDXfox_daily_sync
fox_valuationvaluation腾讯 quotefox_daily_sync
fox_industryindustryTDX + STGfox_wide_sync
fox_holderholder东财 + TDXfox_daily_sync
fox_dividenddividend东财fox_daily_sync
fox_lockup_expirylockup东财fox_daily_sync
fox_dragon_tigerdragon_tiger东财 + THSfox_daily_sync
fox_marginmargin东财fox_daily_sync
fox_block_tradeblock_trade东财fox_daily_sync
fox_announcementannouncement巨潮fox_daily_sync
fox_newsnews东财 + THSfox_daily_sync
fox_kline_dailykline_daily_fox腾讯 HTTPfox_kline_sync
fox_index_dailyindex_daily腾讯 HTTPfox_daily_sync
fox_kline_barkline_bar腾讯 HTTPfox_kline_sync
fox_stock_widestock_wideTDX + 腾讯fox_wide_sync
fox_market_dailymarket_dailyJRJ + TDXfox_daily_sync
fox_market_dpytmarket_dpytJRJ dpytfox_daily_sync
fox_industry_index_dailyindustry_index申万官方fox_wide_sync
fox_etf_dailyetf_daily交易所 + 腾讯etf_daily_sync
fox_fund_flow_dailyfund_flowTDX + 东财兜底fund_flow_sync

ODS 辅助层(计算产物,非融合表)​

定义于 backend/app/models/fox_models_ods.py(15 张表),与 DWD 融合表共用 fox_ 前缀,但语义完全不同:

ODS 辅助表产出任务输入源说明
fox_stock_indicators_dailyquant_calcfox_kline_daily技术指标(RSI/KDJ/MACD/BOLL 等)
fox_stock_pattern_dailyquant_calcfox_kline_dailyK线形态 + 结构形态
fox_stock_cyq_dailyquant_calcfox_kline_daily筹码分布
fox_valuation_percentilequant_calctencent_stock_quote_dailyPE-TTM/PB 历史分位
fox_stock_fscorefscore_calcsina_financial_statementPiotroski F-Score
fox_float_share_dailyquant_calcTDX 财务快照流通股本变动
fox_industry_prosperity_dailyquant_calcfox_stock_wide行业景气度
fox_quote_snapshot实时采集腾讯/TDX盘中实时行情快照
fox_fund_flow_dailyfund_flow_syncTDX + 东财日频资金流向
fox_strategy_pick_dailystrategy_auto_runfox_stock_wide + indicators策略选股结果
fox_indicator_signal_poolquant_calcfox_stock_indicators_daily技术指标信号池
fox_stock_f10_infofox_daily_syncF10 概览公司概况缓存
fox_stock_left_dailyquant_calcfox_kline_daily左Bar数据
fox_stock_mispricing_dailyquant_calcfox_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_TYPESODS 辅助层
目标表fox_* DWD 融合表fox_* 计算产物表(fox_models_ods.py)
绕过方向跳过 STG,融合器直拉 API跳过 STG,计算任务直写
原因API 已返回归一化数据,无 RAW 可落输入源单一、计算确定性高,无多源冲突
注册位置registry.py::NO_ODS_TYPESfox_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_sync 18:50 批量刷新
  • STG truncate-reload:事件 STG 表无业务唯一键,DELETE-all + re-INSERT 是唯一安全的去重方式

相关文档​