数据管理架构
概述
数据管理系统采用多源降级 + 分层缓存架构,确保在任何单一数据源故障时仍能正常提供服务。核心入口为 market_data_service/ 包的 get_market_data()(定义在 fetchers.py,包级 __init__.py 再导出)。
数据源优先级(八步降级)
get_market_data() 用 single-flight 锁包裹整个流程:缓存 miss 时同一 data_type
的并发请求合并为一次获取(前端多组件并发会重复打外部源触发风控)。
| 步骤 | 层级 | 说明 |
|---|---|---|
| 1 | 内存缓存 | key mkt_{data_type};交易时段 TTL 10~30s(按数据变化频率分级),休市 TTL_IDLE=3600s |
| 2 | 非交易日短路 | 周末/节假日行情冻结:直接复用持久层「当日优先、无则最近一日」的快照,不打外部源(历史 bug:曾因要求快照日期 == 今天而周末永远不命中,外部源挂起时把请求拖过网关超时) |
| 3 | 休市优先 DB 当日 | 快照 _snapshot_date 必须等于今天;跳过样本估算快照(sample_based),distribution 补分桶、fund_flow 补资金流字段 |
| 4 | 请求去重 | 休市时若今日已获取过该类型,直接读快照(_request_log,见下节) |
| 5 | 实时源编排 | asyncio.wait_for(_orchestrate_live_sources, timeout=_LIVE_FETCH_BUDGET),总预算超时即落步骤 8;交易时段 TDX → 腾讯,非交易时段的 sentiment/distribution 改为 腾讯 → TDX(腾讯样本估算结果会被 TDX 全量口径顶替) |
| 6 | 东财降级 | 熔断器 + 重试预算双重保护;fund_flow 遇东财限流时改走新浪备用源(source_used="sina_backup") |
| 7 | 持久化 + 缓存 | 有效数据先校正口径(涨跌停权威口径、资金流与前一成交额对齐),再落 fox 持久层并写缓存;样本估算数据不落库,避免休市回退到 ~50 只样本 |
| 8 | DB 历史兜底 | fox 持久层历史数据(同样跳过样本快照),无数据则返回空结构 |
| 数据源 | 协议 | 定位 | 封IP风险 |
|---|---|---|---|
| TDX 通达信(自实现协议) | TCP | 交易时段主源:指数/涨跌统计/成交额实时 | 无 |
| 腾讯财经 | HTTP | 非交易时段 sentiment/distribution 主源,全局第一备胎 | 无 |
| 东方财富 | HTTP | 独有数据(涨停池/人气榜/研报/概念命中)+ 降级兜底 | 高 |
| 新浪 / 百度 | HTTP | 不同域名/风控面,仅特定场景备用(fund_flow、K线 MA) | 低 |
核心文件结构
backend/app/
├── core/
│ ├── exceptions.py # 结构化异常体系
│ ├── cache_manager.py # 内存缓存管理器
│ ├── retry.py # 通用重试装饰器(retry_on_failure)
│ └── database.py # SQLAlchemy 会话管理 + 自动建表 + 增量迁移
├── services/
│ ├── market_data_service/ # 市场数据统一入口(八步降级,包化)
│ │ ├── __init__.py # 再导出 get_market_data(向后兼容)
│ │ ├── fetchers.py # get_market_data 主流程
│ │ ├── infra.py # 连接池/熔断器/去重/TTL/交易时间判定
│ │ ├── intraday_cache.py # 分钟K线轮询缓存(周期自适应 TTL)
│ │ ├── builders.py # 返回结构组装
│ │ ├── persistence.py # 快照落库
│ │ ├── providers_tdx.py # TDX 取数
│ │ ├── providers_tencent.py # 腾讯取数(含全市场涨幅榜 L2 链路)
│ │ └── providers_eastmoney.py # 东财取数(含新浪备用源)
│ ├── data_manager/ # 数据管理器包(门面模式,20+ 子模块)
│ │ ├── __init__.py # DataManager 门面类(170 个 staticmethod 委托)
│ │ ├── downloaders/ # 下载器子包(DOWNLOADER_MAP 59 项)
│ │ │ ├── __init__.py # DOWNLOADER_MAP 汇总 + 各下载器再导出
│ │ │ ├── market.py # K线/行情/估值/资金流
│ │ │ ├── finance.py # 财报三表/财务指标
│ │ │ ├── events.py # 龙虎榜/大宗/两融/解禁/分红
│ │ │ ├── holders.py # 股东户数/十大股东/高管持股/机构调研
│ │ │ ├── news.py # 公告/新闻/研报
│ │ │ ├── cninfo.py # 巨潮专题(披露计划/债券/打新等 15 类)
│ │ │ ├── ths_f10.py # 同花顺 F10 备份源(12 模块)
│ │ │ ├── jrj.py / sina_batch.py / sw_industry.py / em_batch.py
│ │ │ ├── history_kline.py # K线历史数据补全
│ │ │ ├── sina_stmt_batch.py # 新浪财报批量下载
│ │ │ ├── tdx_gpcw.py # 通达信专业财务全字段同步
│ │ │ └── em_wholemarket.py / em_fin_wholemarket.py / base.py
│ │ ├── queries.py # 数据统计(50表)/完整性(26维度)/新鲜度(52维度)/清理
│ │ ├── batch_ops.py # 批量下载编排(DOWNLOADER_MAP 驱动)
│ │ ├── portfolio.py # 自选股导入与管理
│ │ ├── stock_base.py # 全A股(读取走 fox_stock_wide;导入/同步为废弃壳)
│ │ ├── stock_f10.py # F10 公司概况(东财 pc_hsf10)
│ │ ├── tdx_sync.py # 通达信 TCP 同步(板块/财务/除权除息)
│ │ ├── tencent_market.py # 腾讯全市场行情落库/读取(融合层三级读源)
│ │ ├── index_data.py / forex_data.py / future_data.py / us_stock_data.py
│ │ ├── providers_tencent_{stock,index,future,forex,us_stock}.py # 腾讯分域取数
│ │ ├── sse_data.py / szse_data.py / sse_stock_info.py / szse_stock_info.py
│ │ ├── exchange_fusion.py / sse_validator.py # 交易所数据融合与校验
│ │ ├── market_signals.py / sector_snapshot.py / block_kline.py / block_stock.py
│ │ ├── config.py # 配置常量(STOCK_BASE_FIELDS 等)
│ │ ├── config_mgmt.py # 运行时配置读写
│ │ └── sync_log.py # 同步日志 CRUD
│ └── data_sources/
│ ├── tdx_client.py # 通达信 TCP 主客户端(自实现协议,零第三方依赖,26 个 get_* 函数)
│ ├── tdx/ # 自包含协议实现(codec/connection/commands/client/f10)
│ ├── tdx_client_compat.py # 兼容层(re-export tdx_client;原名 mootdx_client.py,已于 2026-10-06 改名)
│ ├── tdx_client_backup_pool.py # 降级备用(独立服务器池的第二条连接路径)
│ ├── tencent_client.py # 腾讯财经 HTTP 客户端
│ ├── em_client.py # 东财限流入口(em_get)
│ ├── em_paginated.py # 东财分页拉取 + push2轮转(主机池唯一定义源)
│ ├── baostock_client.py # 历史兼容壳(转发东财 push2his,已无调用方)
│ └── real_data_source.py # 组合数据源实现
├── api/
│ ├── market.py # /api/v1/market/* 路由
│ ├── datamgr.py # /api/v1/datamgr/* 路由(含 38 种按股下载类型 + 新鲜度检查)
│ ├── stocks.py # /api/v1/stocks/* 路由
│ └── … # 共 72 个路由模块(任务中心/AI 分析/量化选股/跨市场品种/鉴权设置),
│ # 由 discover_routers() 自动发现注册,速查见 docs/docs/reference/api-routes.md
└── models/
└── db_models.py # 126 张表 ORM 定义(含 FoxKlineDaily/FinancialIndicator 等;
# 全库共 215 张,余量在 fox_models*/stg_models/sse*/szse*/task_models)
DataManager 门面模式
DataManager 类位于 data_manager/__init__.py,采用**门面模式(Facade)**将所有子模块方法委托为 staticmethod:
- 数据下载(
DOWNLOADER_MAP共 59 个下载器;其中 38 项挂在下表/download的types上(与_DOWNLOAD_TYPE_MAP一一对应),其余 21 项为市场级/专题下载器,只走各自专用端点:巨潮专题 13(cninfo_*:股票列表/定期报告/披露预约/分红明细/股东大会/股本变动/融资融券/大宗统计与明细/债券/互联互通及活跃榜/打新)、金融界 4(jrj_history/jrj_market/jrj_timeline/jrj_yzzdt)、东财 2(em_finance_main/em_org_basic)、申万 1(sw_industry)、全市场财务指标 1(fin_indicator_all)):download_kline,download_quote,download_valuation,download_fund_flow,download_reports,download_dragon_tiger,download_margin,download_holder_num,download_lockup,download_block_trade,download_dividend,download_concept_blocks,download_announcements,download_news,download_financial_statements,download_financial_indicator,download_issue_info,download_holding_org,download_top_holders,download_earnings_forecast,download_earnings_express,download_executive_holding,download_org_survey,download_northbound,download_fund_holding,download_minute_kline+ 同花顺 F10 12 模块(download_ths_*) - 批量操作:
download_stock_all,download_batch - 数据查询:
get_data_stats(50 张表),get_stock_data_status(26 类维度),check_data_freshness(52 类关键维度) - 数据清理:
cleanup_old_data,cleanup_all_old - 配置/日志:
get_config,set_config,get_all_config,create_sync_log,update_sync_log,get_sync_logs - 自选股:
import_portfolio_excel,get_portfolio_stocks - 全A股:
get_stock_base_list(读fox_stock_wide);import_stock_base/sync_stock_base_from_eastmoney等导入同步方法为废弃壳(stock_base_info已移除),恒返回{"status": "deprecated"} - F10:
sync_f10_company_profile,get_f10_info - TDX:
sync_block_info,sync_market_stat,sync_xdxr_info(sync_stock_list_from_tdx/sync_finance_from_tdx同样是stock_base_info废弃壳)
所有单股下载器通过 downloaders/__init__.py 的 DOWNLOADER_MAP dict(59 项)映射到字符串键,供 batch_ops 和 API 层动态调用。
数据类型
get_market_data(data_type) 支持以下数据类型(TTL 默认值来自 core/config.py,可经环境变量覆盖):
| data_type | 说明 | 交易时段TTL | 休市TTL |
|---|---|---|---|
sentiment | 市场情绪(指数+涨跌分布+情绪分数) | 10s | 3600s |
distribution | 涨跌家数分布 | 10s | 3600s |
fund_flow | 资金流向(成交额+主力净流入) | 15s | 3600s |
industry_rank | 行业板块涨跌排名 | 30s | 3600s |
hot_stocks | 热门股票 TOP N | 10s | 3600s |
sector_treemap | 板块热力图数据 | 15s | 3600s |
请求去重机制
同一 data_type 同一天仅调用一次外部 API,后续请求直接读 DB(实现位于 infra.py):
_request_log: dict[str, str] = {} # {data_type: "2026-09-21"}
def already_fetched_today(data_type: str) -> bool:
return _request_log.get(data_type) == today_str()
def mark_fetched(data_type: str): ... # 仅在“有效数据成功取得”后调用
注意:去重仅在休市时生效(交易时段数据持续变化,必须实时拉取);
实时编排超时的路径不写 mark_fetched、不落快照,语义与"获取失败"一致,下次请求会重试。
并发控制
- 线程池:
ThreadPoolExecutor(max_workers=8, thread_name_prefix="mkt")执行同步 TDX/DB 调用 - asyncio.gather:并发获取互不依赖的数据(如指数 + 涨跌分布)
- 共享连接池:
httpx.AsyncClient单例,max_connections=20
# 并发获取示例(sentiment 类型)
indices_task = loop.run_in_executor(_executor, _tdx_indices)
breadth_task = loop.run_in_executor(_executor, _db_market_breadth)
indices, breadth = await asyncio.gather(indices_task, breadth_task)
数据持久化
成功获取的数据通过 fox 门面写入持久层(market_snapshot 通用快照表已于 2026-08-22 退役迁移至 fox 层):
- fund_flow →
fox_market_daily(盘后定型权威数据) - 其余 5 类(sentiment/distribution/industry_rank/hot_stocks/sector_treemap)→
fox_market_snapshot(date+data_type 双主键,当日每类一行,幂等)
-- fox_market_snapshot 结构(替代已退役的 market_snapshot)
CREATE TABLE fox_market_snapshot (
date VARCHAR(10) NOT NULL, -- 数据日期 YYYY-MM-DD
data_type VARCHAR(32) NOT NULL, -- sentiment/distribution/industry_rank/hot_stocks/sector_treemap
data_json TEXT,
source VARCHAR(32) DEFAULT 'tdx', -- 数据来源
quality VARCHAR(10) NOT NULL DEFAULT 'ok',
updated_at DATETIME,
PRIMARY KEY (date, data_type),
INDEX idx_fms_type_date (data_type, date DESC)
);
自动建表机制
应用启动时自动创建缺失的数据库表(database.py 的 ensure_tables()):
def ensure_tables():
"""幂等创建缺失表,已存在的表不受影响。"""
if engine is None:
return
from ..models import db_models # 确保 Base.metadata 包含全部表定义
Base.metadata.create_all(bind=engine, checkfirst=True)
- 在
main.py的lifespan启动钩子中调用 - 解决
init_db.sql未覆盖新增表(如fox_stock_master、app_portfolio_stocks)的问题 - 幂等安全:
checkfirst=True确保已有表不会被重建或覆盖
全A股数据管理(fox 融合层)
早期「文件导入 / 东财在线同步 / TDX 同步」三种全A股入库方式均以 stock_base_info(149 字段
TDX 全量导出)与精简表 stock_info 为目标表,这两张表已于 2026-08-22 移除。
POST /stock-base/import|sync|sync-tencent|sync-tdx|sync-full|backfill|fetch-staging|promote|
backfill-stock-info、POST /finance/sync-tdx 与 GET /stock-base/staging 现为兼容壳
(恒返回 {"status": "deprecated", ...},不写任何表)。
全A股数据现在完全由 fox 融合层产出,无需人工导入:
| 层次 | 表 | 产出方 | 说明 |
|---|---|---|---|
| 主数据 | fox_stock_master | fox_wide_sync(TDX 证券列表 + 巨潮) | 代码/名称/股本/上市日期/股东人数 |
| 行情宽表 | fox_stock_wide | fox_wide_sync(TDX + 腾讯批量行情) | 价格/涨跌幅/换手/量比/PE/PB/市值/主力净额/行业 |
| 行业归属 | fox_industry | fox_wide_sync(TDX + 申万 STG) | 腾讯三级 industry + 申万 sw_l1/sw_l2 代码 |
- 触发方式:定时任务
fox_wide_sync(交易日 16:20),或手动POST /datamgr/fox/sync(全市场模式需持有三域锁,冲突返回 409;详见 datamgr-api) - 读取入口:
GET /datamgr/stock-base/list(数据源即fox_stock_wide,保留旧分页契约与单位口径)
仍在使用 TDX 的其余同步接口:
| 方法 | API | 说明 |
|---|---|---|
| 板块同步 | POST /block/sync | 行业/概念/风格板块 → tdx_block_info(按 block_type 先删后写) |
| 市场统计 | GET /market/stat | 涨跌家数/涨停跌停/总市值(TDX TCP 实时,失败返回 503) |
| 除权除息 | GET /stock/{code}/xdxr | 分红送转/股本变动历史 |
| 财务同步 | POST /finance/sync-tdx | 已废弃(stock_base_info 已移除;财务改走 fox_finance_indicator) |
操作日志记入 app_sync_log(任务中心的 app_task_run 为准,app_sync_log 为兼容层)。
RAW → STG → DWD 三层数据管线
fox 融合层采用 RAW → STG → DWD 三层管线,将数据采集、清洗归一、交叉融合三个阶段解耦:
RAW(30+ 张 em_/tdx_/ths_/sina_/cninfo_/jrj_ 表)
│ ← 各下载器逐源落库,只管写自己的 RAW 表
▼
STG(14 张 stg_* 表)
│ ← StgLoaderSpec 声明式注册,run_stg_loader 统一驱动
│ ← 归一化 / 去重 / 血缘追踪 / 质量校验
▼
DWD(22 张 fox_* 表)
│ ← FoxFuserSpec 声明式注册,run_fuser 统一驱动
│ ← 多源交叉验证 / 融合 / 口径对齐
▼
消费端(API / 前端 / 量化指标 / 评分)
STG 层的角色
STG(Staging,清洗贴源层)位于 RAW 与 DWD 之间,承担 ODS + 数据治理 双重职责:
| 职责 | 说明 |
|---|---|
| 归一化 | 不同 RAW 源的字段名/日期格式/数值单位统一为 STG 标准列 |
| 血缘追踪 | 每行 raw_source + source_json 记录来自哪些 RAW 表 |
| 质量校验 | 空值率/数值范围检查 → quality 列(ok / warning) |
| 解耦采集与融合 | RAW 下载器只管落库,STG loader 负责清洗,DWD 融合器只消费 STG |
| 幂等 UPSERT | 重复执行覆盖,不产生重复数据(业务键 ON DUPLICATE KEY UPDATE) |
声明式注册
两类 Spec 采用相同模式——数据类型声明 + 三函数(read_raw / list_codes / build_rows):
| 注册表 | 位置 | 条目数 | 执行入口 |
|---|---|---|---|
STG_LOADERS | stg_loader/__init__.py | 14 | sync_stg_all() / sync_stg_by_types() |
FOX_FUSERS | writer/registry.py | 22 | sync_all() / sync_by_types() |
STG 与 DWD 类型映射
14 个 STG 类型与 22 个 DWD 融合类型的对应关系:
| STG 类型 | DWD 融合类型 | 说明 |
|---|---|---|
| stock_master | stock_master | 主数据:代码/名称/股本/上市日期 |
| company_profile | company_profile | F10 公司概况 |
| finance_indicator | finance_indicator | 财务指标:营收/净利/ROE/EPS |
| valuation | valuation | 估值:PE/PB/市值 |
| industry | industry | 行业归属:腾讯三级 + 申万 L1/L2 |
| holder | holder | 股东户数 |
| dividend | dividend | 分红送转 |
| lockup | lockup | 解禁 |
| dragon_tiger | dragon_tiger | 龙虎榜 |
| margin | margin | 两融 |
| block_trade | block_trade | 大宗交易 |
| announcement | announcement | 公告 |
| news | news | 新闻/研报 |
| market_daily | market_daily | 大盘日统计 |
8 类绕过 STG 的融合类型(NO_ODS_TYPES)
以下 8 类 DWD 融合没有 STG 层,融合器直接从 API 取数:
| 类型 | 原因 |
|---|---|
| kline_daily_fox | 腾讯 fqkline 日K:API 直拉,KlineStore 为并行支路非 ODS |
| kline_bar | 腾讯 fqkline 周/月/季/年K:API 直拉 |
| index_daily | 腾讯指数日K:API 直拉,指数清单为静态配置 |
| stock_wide | 个股宽表:多头实时快照按需现算 |
| market_dpyt | 大盘云图:JRJ dpyt 行业树 API 直拉 |
| industry_index | 申万行业指数日线:官方 trend 接口直拉,无日期参数 |
| macro_indicator | 宏观指标最新一期:多源 API 直拉,月/季频快照 |
| macro_release | 宏观发布日程:规则推算 + API 直拉 |
绕过理由统一登记在 registry.py::NO_ODS_TYPES,由 tests/test_layer_boundaries.py 锁住「未登记理由的融合类型不允许绕过 STG」。
两阶段执行流
sync_all() 按 Phase 1 → Phase 2 顺序执行:
# Phase 1: STG 层(RAW → STG,14 类)
sync_stg_all(codes, batch_id)
# Phase 2: DWD 层(STG → fox_*,22 类)
for data_type in _FOX_ALL_TYPES:
run_fuser(FOX_FUSERS[data_type], ...)
- Phase 1 先于 Phase 2:DWD 融合器读 STG 表,必须先完成清洗归一
- Phase 2 内部按依赖排序:stock_master → industry → stock_wide(宽表读主数据与行业)
- 三域独立任务:
fox_kline_sync(15:40) /fox_daily_sync(16:00) /fox_wide_sync(16:20) 各自覆盖部分类型,类型不相交、可并行
STG 清理机制
STG 层时序表已纳入 data_cleanup 定时任务的滚动清理范围(2026-09-29):
| 类别 | 表数量 | 清理策略 |
|---|---|---|
| 时序表 | 11 张 | 按日期列滚动清理,默认保留 3 年 |
| 维度表 | 3 张 | 不纳入清理(无日期列,每次全量覆盖写) |
纳入清理的 11 张 STG 时序表:
| 表名 | 日期列 | 说明 |
|---|---|---|
| stg_valuation | date | 估值快照 |
| stg_holder | date | 股东户数 |
| stg_dividend | ex_date | 分红送转 |
| stg_lockup | date | 解禁 |
| stg_dragon_tiger | date | 龙虎榜 |
| stg_margin | date | 两融 |
| stg_block_trade | date | 大宗交易 |
| stg_announcement | pub_date | 公告 |
| stg_news | pub_time | 新闻/研报 |
| stg_market_daily | date | 大盘日统计 |
| stg_finance_indicator | updated_at | 财务指标(按清洗时间) |
不纳入清理的 3 张 STG 维度表:
| 表名 | 原因 |
|---|---|
| stg_stock_master | 无日期列,每次全量覆盖写 |
| stg_company_profile | 无日期列,每次全量覆盖写 |
| stg_industry | 无日期列,每次全量覆盖写 |
清理参数:
- 执行时机:每周一 04:30(
data_cleanup定时任务) - 默认保留:1095 天(3 年)
- 最小保留:30 天(防止误删近期数据)
- 配置方式:SysConfig
data.cleanup.retention_days
设计理由:
- STG 表数据可从 RAW 层重建,但大表(margin/announcement 等百万行级)长期不清理会浪费存储
- 维度表每次全量覆盖写,无历史累积,不需要清理
- STG 表未加入
PROTECTED_HISTORY_TABLES(与 DWD 层历史保护表不同),因为 STG 可从 RAW 重建,默认 3 年保留窗口已足够