Skip to Content
数据基石后端模块数据采集模块详细设计

ZLData 数据采集模块 — 详细设计文档

属性内容
文档编号ZLData-DC-DD-001
配套文档ZLData_PRD_v2.md v2.1
日期2026-04-03
状态Draft — 待评审

目录

  1. 模块概述
  2. 数据源设计
  3. 双源冗余架构
  4. 行情数据采集
  5. 财务数据采集
  6. 宏观经济数据采集
  7. 利率数据采集
  8. 新闻事件采集与情感分析
  9. 机构研报采集
  10. 财经日历与催化剂发现
  11. 数据标准化与质量保障
  12. 数据入库设计
  13. 性能优化策略
  14. 采集调度策略
  15. 配置管理
  16. 异常处理与监控

1. 模块概述

1.1 定位与目标

数据采集模块(Data Collection)是 ZLData AI 量化交易平台的 数据基础设施层,承担着从多源异构数据源获取行情、财务、宏观、新闻和研报数据,并经过清洗、标准化、去重后写入中心数据库的核心职责。该模块是整个平台的数据基石——量化分析、策略回测、风险评估和智能体决策的上游输入完全依赖本模块所提供的数据质量与时效性。

设计目标包括:

  • 数据完整性:覆盖技术分析(行情)、基本面分析(财务)、宏观环境(指标)和另类数据(新闻/研报/事件)四类数据,构建完整的 AI 特征库。
  • 高可用性:通过双数据源冗余架构,确保单一数据源故障时数据采集不中断,数据可用性达到 99.99%。
  • 时效性保障:日内数据延迟 < 5 秒,日频数据在收盘后 30 分钟内完成更新。
  • 幂等与可恢复:任何采集任务可安全重复执行,断网或崩溃后从断点续传,无需删表重来。
  • 生产级质量:包含完整的异常处理、数据质量检查、限速策略和监控告警机制。

1.2 数据分类体系

1.3 核心数据表清单

序号表名说明数据量级更新频率
1trade_stock_daily日K线行情~5000只 × 3000日 ≈ 1500万行日频
2trade_stock_financial季度财务指标~5000只 × 40期 ≈ 20万行季频
3trade_macro_indicator月度宏观指标~120行(10年)月频
4trade_rate_daily日频利率指标~2500行(10年)日频
5trade_stock_news个股新闻事件~50万行(年增量)小时级
6trade_report_consensus机构研报一致性~10万行周频
7trade_calendar_event财经日历事件~2000行(年增量)日频

2. 数据源设计

2.1 数据源矩阵

2.2 数据源能力对比

能力维度QMT (xtquant)akshare定向爬虫LLM搜索
日K线数据✅ 实时/历史,前/后复权✅ 历史日/周/月线
分钟线/Tick✅ 1min/5min/Tick
财务报表✅ 5张报表(Balance/Income/CashFlow/PershareIndex/Capital)✅ 新浪三大报表(利润/资产负债/现金流)
宏观指标✅ CPI/PPI/PMI/M2/LPR等
利率数据✅ 国债收益率
新闻资讯✅ 东方财富个股新闻✅ 财联社/同花顺
机构研报✅ 东方财富评级 + 同花顺预测✅ 券商官网
财经日历✅ 百度财经日历
催化剂事件✅ LLM搜索发现
板块/行业✅ 申万一级/二级行业分类✅ 指数成分
数据延迟< 1秒(本地缓存)1-3秒(网络请求)5-30秒10-60秒
频率限制无(本地数据)中(需限速防封)低(需反爬)按API额度
稳定性★★★★★ (本地服务)★★★☆☆ (依赖网站)★★☆☆☆ (易被封)★★★★☆ (需付费API)

2.3 akshare API 映射表

数据类型akshare 函数数据来源说明
日K线ak.stock_zh_a_hist()东方财富支持 daily/weekly/monthly,前/后复权
财务报表ak.stock_financial_report_sina()新浪财经利润表/资产负债表/现金流量表
宏观CPIak.macro_china_cpi_monthly()国家统计局月度同比/环比/累计
宏观PPIak.macro_china_ppi_monthly()国家统计局月度同比/环比
宏观PMIak.macro_china_pmi()国家统计局制造业/非制造业PMI
宏观M2ak.macro_china_money_supply()人民银行M0/M1/M2同比
社融数据ak.macro_china_society_financing()人民银行社融规模增量
LPR利率ak.macro_china_lpr()人民银行1年期/5年期
国债收益率ak.bond_china_yield()中国债券网日频,多期限
个股新闻ak.stock_news_em()东方财富按股票代码获取
机构评级ak.stock_institute_recommend_detail()东方财富券商/评级/目标价
盈利预期ak.stock_profit_forecast_ths()同花顺EPS/净利润一致预期
财经日历ak.news_economic_baidu()百度财经按日期获取经济事件

3. 双源冗余架构

3.1 架构总览

3.2 QMT 主数据源

QMT(迅投极速策略交易系统)通过 xtquant Python SDK 提供本地化数据服务。核心优势在于数据存储在本地缓存中,读取速度极快,且支持全量 A 股约 5000 只股票的完整历史数据。

连接流程

from xtquant import xtdata # 1. 连接本地QMT服务 connect_result = xtdata.connect() # 注意:需先启动 miniQMT 客户端 # 2. 获取股票列表 all_codes = xtdata.get_stock_list_in_sector('沪深A股') # 3. 获取行业分类 sw1_sectors = [s for s in xtdata.get_sector_list() if str(s).upper().startswith('SW1') and '加权' not in str(s)]

关键 API

API功能参数
xtdata.connect()连接QMT服务
xtdata.get_stock_list_in_sector(sector)获取板块成分股板块名如’沪深A股’
xtdata.download_history_data()下载历史行情到本地缓存stock_code, period, start_time
xtdata.get_market_data_ex()获取行情数据stock_list, period, dividend_type
xtdata.download_financial_data2()异步下载财务数据(推荐)stock_list, table_list, callback
xtdata.get_financial_data()获取财务数据stock_list, table_list, report_type
xtdata.get_sector_list()获取板块列表
xtdata.get_stock_list_in_sector()获取板块成分板块名
xtdata.get_instrument_detail()获取股票详情(名称等)stock_code

3.3 akshare 备用数据源

akshare 是一个免费开源的 Python 金融数据库,通过模拟浏览器访问新浪财经、东方财富、国家统计局等官网获取数据。无需 Token 或注册即可使用。

与 QMT 的数据差异处理

维度QMTakshare统一策略
股票代码格式600519.SH600519内部统一为 600519.SH
日期格式20260101(字符串) / 时间戳2026-01-01(字符串)统一转换为 DATE 类型
财务报表字段s_fa_eps_basic, tot_assets中文列名”基本每股收益”等各自映射到标准字段名
K线列名open, close, volume开盘, 收盘, 成交量akshare端做列名重映射
数据精度高(专业数据源)中(网页解析)标记 data_source 字段区分

3.4 数据源切换判定条件

判定条件动作
QMT 连接超时 (>30s)切换 akshare
QMT 返回空数据切换 akshare
QMT 数据量异常(如日K线 < 10条且历史 > 1年)切换 akshare 并记录告警
akshare 请求失败 (HTTP 5xx / 超时)重试 3 次后标记失败
双源均不可用告警 P1,30分钟后重试

4. 行情数据采集

4.1 前复权处理

为什么必须复权?

A股市场中,上市公司会进行分红、送股、配股等操作。以贵州茅台为例:若每10股送10股(高送转),次日股价理论上从100元变为50元。如果不做复权处理,AI 模型会将此视为暴跌(-50%),触发错误的卖出信号。

复权方式选择

复权方式特点适用场景
前复权以当前价格为基准,历史价格等比例缩小✅ 回测与交易(保持当前价格真实)
后复权以首日价格为基准,后续价格等比例放大❌ 不推荐(当前价格失真)
不复权原始价格,含分红除权缺口❌ 不推荐(指标计算会失真)

QMT 调用

data = xtdata.get_market_data_ex( stock_list=[stock_code], period='1d', dividend_type='front' # 前复权 )

akshare 调用

df = ak.stock_zh_a_hist( symbol='600519', period="daily", start_date='20240101', end_date='20251231', adjust="qfq" # 前复权 )

4.2 换手率计算

QMT 的原始行情数据不包含换手率字段,需根据流通股本自行计算。换手率是衡量资金活跃程度的关键指标,换手率突增往往意味着主力进场或重大消息刺激。

计算公式

TurnoverRate = (Volume × 100 / FloatingShares) × 100%

其中:

  • Volume:QMT 返回的成交量,单位为 (1手 = 100股),需乘以 100 转为股
  • FloatingShares:流通股本,单位为股,通过 xtdata 获取

代码实现

def _get_float_shares(stock_code: str) -> float: """从QMT获取流通股本""" detail = xtdata.get_instrument_detail(stock_code) return float(detail.get('TotalShares', 0)) if detail else 0 float_shares = _get_float_shares(stock_code) for idx, row in df.iterrows(): vol = int(row['volume']) vol_shares = vol * 100 # 手 → 股 turnover = round(vol_shares / float_shares * 100, 4) \ if float_shares > 0 and vol > 0 else None

4.3 增量更新策略

全量A 股有约 300 多万条日线数据(5000只 × ~600个交易日),每天全量重下不现实。采用 增量更新 策略:

批量查询最新日期的 SQL

SELECT stock_code, MAX(trade_date) AS max_date FROM trade_stock_daily GROUP BY stock_code

返回 {stock_code: max_date} 的映射,避免对每只股票单独查询(从 5000 次 SQL 减少到 1 次)。

4.4 akshare 数据格式标准化

akshare 返回中文列名,需统一映射为英文标准字段名:

COL_MAP = { '日期': 'date', '开盘': 'open', '收盘': 'close', '最高': 'high', '最低': 'low', '成交量': 'volume', '成交额': 'amount', '振幅': 'amplitude', '涨跌幅': 'pct_chg', '涨跌额': 'change', '换手率': 'turnover', } df = df.rename(columns=COL_MAP) df['date'] = pd.to_datetime(df['date'])

5. 财务数据采集

5.1 QMT 财务报表体系

QMT 通过 xtdata.get_financial_data() 提供 5 张财务报表数据:

报表表名关键字段
资产负债表Balancetot_assets, tot_liab, total_equity, total_current_assets, total_current_liability, inventories, cash_equivalents
利润表Incomerevenue, net_profit_incl_min_int_inc, operating_revenue, oper_profit, cost_of_goods_sold
现金流量表CashFlownet_cash_flows_oper_act, net_cash_flows_inv_act, net_cash_flows_fnc_act
每股指标表PershareIndexs_fa_eps_basic, s_fa_bps, s_fa_ocfps, s_fa_undistributedps, du_return_on_equity, sales_gross_profit
股本表CapitaltotalShares, totalCapital

5.2 多候选字段匹配

由于 QMT 不同版本可能使用不同的字段名,设计多候选字段查找机制:

def get_field(record: dict, field_names: list, default=None): """ 从记录中获取字段值,支持多个候选字段名。 按优先级依次尝试,返回第一个非空值。 """ for name in field_names: val = record.get(name) if val is not None: return val return default

典型使用场景

# 净利润:优先含少数股东损益版 net_profit = get_field(inc, ['net_profit_incl_min_int_inc', 'net_profit_excl_min_int_inc']) # 总资产 total_assets = get_field(bal, ['tot_assets']) # ROE:优先取PershareIndex预计算值,没有则手动计算 roe = get_field(ps, ['du_return_on_equity', 'equity_roe', 'net_roe']) if roe is None and net_profit and total_equity: roe = safe_divide(net_profit, total_equity, pct=True)

5.3 衍生指标计算

部分指标在原始数据中可能缺失,需从基础报表数据计算:

def safe_divide(a, b, pct=False): """安全除法,b为0时返回None。pct=True时结果乘以100。""" if a is None or b is None: return None a, b = float(a), float(b) if b == 0: return None result = a / b return round(result * 100, 4) if pct else round(result, 4)

指标计算映射

衍生指标计算公式数据来源优先级
ROE净利润 / 净资产 × 100Income + Balance先取 PershareIndex,再计算
ROA净利润 / 总资产 × 100Income + Balance始终计算
毛利率(营收 - 营业成本) / 营收 × 100Income先取 PershareIndex,再计算
净利率净利润 / 营业收入 × 100Income始终计算
资产负债率总负债 / 总资产 × 100Balance始终计算
流动比率流动资产 / 流动负债Balance始终计算
速动比率(流动资产 - 存货) / 流动负债Balance始终计算
总资产周转率营业收入 / 总资产Income + Balance始终计算
OCF/营收经营现金流 / 营业收入 × 100CashFlow + Income始终计算
OCF/净利润经营现金流 / 净利润CashFlow + Income始终计算

5.4 净利润同比增长率

取年报数据(报告期以 1231 结尾),计算最近两年的同比增长率:

def calc_netprofit_yoy(records: list) -> float | None: """用年报数据计算最新一期净利润同比增长率。""" annual = [r for r in records if str(r.get('end_date', '')).endswith('1231')] if len(annual) < 2: return None annual = sorted(annual, key=lambda x: x['end_date']) profits = [float(r['net_profit']) for r in annual if r.get('net_profit') is not None] if len(profits) < 2: return None s = pd.Series(profits) return round(s.pct_change().iloc[-1] * 100, 2)

5.5 批量下载策略

全量A 股财务数据量大,采用 每批 50 只股票 的批量下载策略:

性能对比

方式API 调用次数预估耗时
逐只下载~5000次~3小时
50只/批~100次~15分钟

5.6 财务数据采集完整流程


6. 宏观经济数据采集

6.1 采集指标清单

宏观经济数据为 AI 模型提供 外部环境特征,帮助模型理解为什么同样的财务表现在不同年份股价表现天差地别。例如:特斯拉 2021 年净利润 55 亿美元(10年期美债 1.5%,PE > 200x,股价冲上 400 美元);2022 年净利润翻倍至 126 亿美元,但美债飙升至 4.3%,PE 跌至 30x,股价暴跌 65%。

指标akshare 函数更新频率重要性
CPI 同比ak.macro_china_cpi_monthly()月频★★★★★
PPI 同比ak.macro_china_ppi_monthly()月频★★★★★
PMIak.macro_china_pmi()月频★★★★★
M2 同比增速ak.macro_china_money_supply()月频★★★★☆
社融规模增量ak.macro_china_society_financing()月频★★★★☆
LPR 1年期ak.macro_china_lpr()月频★★★★☆
LPR 5年期ak.macro_china_lpr()月频★★★★☆

6.2 日期清洗与合并

宏观数据来自不同官方机构,日期格式五花八门。akshare 返回的日期可能是 2026年01月2026012026-01 等多种格式,需要统一清洗为 DATE 类型。

由于各指标公布时间不同(有的月初,有的月中),采用 Outer Join 合并策略,确保即使某指标本月尚未公布,其他已公布的指标也能正常写入。

6.3 COALESCE 保护模式

COALESCE 是 SQL 的空值处理函数,返回第一个非 NULL 的值。在宏观数据更新时,使用 COALESCE 确保已有有效数据不会被接口故障返回的 NULL 值覆盖。

INSERT INTO trade_macro_indicator (indicator_date, cpi_yoy, ppi_yoy, pmi, data_source) VALUES ('2026-02-28', 2.1, -1.3, 50.8, 'akshare') ON DUPLICATE KEY UPDATE cpi_yoy = COALESCE(VALUES(cpi_yoy), cpi_yoy), ppi_yoy = COALESCE(VALUES(ppi_yoy), ppi_yoy), pmi = COALESCE(VALUES(pmi), pmi);

保护逻辑:如果 akshare 某次接口调用失败返回了 NULL,COALESCE 会保留数据库中已有的历史值,确保数据采集具备容错能力。


7. 利率数据采集

7.1 采集内容

利率数据为 AI 模型提供 流动性环境特征。10年期国债收益率被誉为”全球资产定价之锚”,它的上涨意味着市场估值承压,下跌则意味着流动性宽松。

指标akshare 函数更新频率说明
中国10年期国债收益率ak.bond_china_yield()日频资产定价之锚(A股端)
美国10年期国债收益率ak.bond_china_yield()日频北向资金流向指标
中美利差计算值日频cn_bond_10y - us_bond_10y,利差倒挂信号

7.2 采集流程

df = ak.bond_china_yield(symbol="中国国债收益率") # 提取10年期收益率列 cn_10y = df[['date', '中国国债收益率10年']] cn_10y.columns = ['rate_date', 'cn_bond_10y']

写入 trade_rate_daily 表,使用 ON DUPLICATE KEY UPDATE 确保幂等。


8. 新闻事件采集与情感分析

8.1 两层情感分析策略

全量A 股每天产生海量新闻,全部用大模型处理成本极高。采用 关键词粗筛 + LLM 精读 双层策略:

关键词库

类别关键词
利好涨停、超预期、增长、中标、突破、新高、回购、增持、大单、主力
利空大跌、减持、亏损、处罚、暴雷、退市、违规、质押、商誉、减持

8.2 标题去重

同一条新闻(如”今日白酒板块表现强劲”)可能出现在多只股票的公告栏里。入库前检查标题是否已存在,存在则跳过,避免数据库膨胀。

8.3 并行采集与限速

from concurrent.futures import ThreadPoolExecutor, as_completed import threading NUM_WORKERS = 8 # 8个线程并行 _print_lock = threading.Lock() # 线程锁,防止打印乱序 # 增量跳过:查询今日已采集股票 skip_sql = """ SELECT DISTINCT stock_code FROM trade_stock_news WHERE DATE(created_at) = CURDATE() """

限速原则:8 线程并行可以提速,但线程过多容易被目标网站封禁 IP。建议搭配请求间隔(每线程每次请求间隔 0.5-1 秒)。


9. 机构研报采集

9.1 变更检测去重

同一家券商会在”行业晨报”、“周报”、“深度报告”中反复提到同一个观点。噪音数据会干扰 AI 模型。

去重策略:只记录变化。查询数据库获取每家券商对该股票的最新观点,只有当新采集到的评级或目标价与数据库记录不同时,才写入新记录。

9.2 量化因子

研报数据可构建两个高价值的 AI 量化因子:

一致预期偏离度因子

Deviation = (Actual_EPS - Consensus_EPS) / Consensus_EPS

逻辑:股价不是由好消息驱动的,而是由比预期更好的消息驱动。如果实际 EPS 远超分析师一致预期,往往迎来超预期上涨。

目标价空间因子

Upside = (Target_Price - Current_Price) / Current_Price

策略:在 AI 选股模型中,给目标空间 > 20% 且近期评级上调的股票增加权重。

9.3 7天跳过策略

研报数据不像行情数据那样秒变,更新频率较低。对每只股票,如果过去 7 天内已采集过,则直接跳过,降低请求频率。

SELECT DISTINCT stock_code FROM trade_report_consensus WHERE created_at >= DATE_SUB(NOW(), INTERVAL 7 DAY)

10. 财经日历与催化剂发现

10.1 财经日历采集

财经日历记录未来重要经济事件的发布时间,为量化模型提供”时间坐标系”。核心价值在于风险规避(重大数据公布前降低仓位)和为 LLM 提供上下文。

A股重点关注事件

事件发布频率影响板块
LPR 利率每月20日房地产、银行
社融/M2每月中旬全市场(流动性)
PMI每月底周期性行业
美联储议息(FOMC)每季度北向资金、外资重仓股
政治局会议4/7/10/12月全市场(政策定调)

10.2 LLM 驱动的催化剂发现

普通财经日历只能抓到定期发布的数据(如 CPI),但很多影响股价的事件(突发会议、行业龙头财报预警、政策变动)没有标准 API。采用 LLM 联网搜索 + 两阶段处理 来发现这些催化剂事件。

10.3 Prompt 管理

生产环境中,Prompt 不应硬编码在代码中,使用 YAML 文件管理:

# prompts.yaml calendar: search_catalysts: | 你是一位顶级的宏观策略分析师。请联网搜索并整理从{start_date} 到 {end_date} 可能对 A 股产生"定调性"或"脉冲性"影响的事件。 你需要重点关注: 1. 决策层动态:如政治局会议(4/7/10/12月)、中央经济工作会议。 2. 产业节点:如华为/特斯拉的发布会、英伟达财报日、全球AI峰会。 3. 市场规则:如MSCI/富时罗素指数调整生效日。 请以 JSON 数组格式返回,每个事件包含 date, title, country, category, importance。

11. 数据标准化与质量保障

11.1 安全类型转换

所有数据采集模块共享的安全转换工具函数:

def safe_float(val) -> float | None: """安全转换为浮点数""" if val is None or str(val).strip() in ('', '--', 'None', 'nan'): return None try: return float(val) except (ValueError, TypeError): return None def safe_divide(a, b, pct=False) -> float | None: """安全除法,b为0时返回None""" if a is None or b is None: return None a, b = float(a), float(b) if b == 0: return None result = a / b return round(result * 100, 4) if pct else round(result, 4)

11.2 日期标准化

不同数据源返回不同的日期格式,统一处理策略:

原始格式来源标准化结果
20260331QMT (m_timetag)2026-03-31 (DATE)
1709251200000QMT (毫秒时间戳)2024-03-01 (DATE)
2026年03月akshare 宏观2026-03-31 (月末日期)
202603akshare 宏观2026-03-31 (月末日期)
2026-03-31akshare 行情2026-03-31 (DATE)

QMT 时间戳处理

def normalize_timetag(ts_val): """将xtquant的m_timetag转换为日期字符串""" if ts_val is None: return None s = str(ts_val).strip() if len(s) == 8 and s.isdigit(): return s try: v = float(s) if v == 0: return None if v > 1e12: v = v / 1000 # 毫秒 → 秒 return datetime.fromtimestamp(v).strftime('%Y%m%d') except (OSError, ValueError, TypeError): return None

11.3 数据质量维度

维度指标检查方法告警阈值
完整性字段非空率COUNT(非空) / COUNT(总)核心字段 > 95%
时效性数据更新延迟MAX(当前时间 - 数据时间)日K线 < 1小时
一致性跨源数据偏差|QMT值 - akshare值| / QMT值偏差 < 1%
准确性异常值检测Z-score > 3 或值为 0标记为待核查
唯一性重复记录数COUNT - COUNT(DISTINCT 唯一键)重复数 = 0

12. 数据入库设计

12.1 幂等写入模式

核心原则:无论执行一次写入操作,还是重复执行一万次,数据库的最终状态完全一致。

-- 幂等写入示例 INSERT INTO trade_stock_daily (stock_code, trade_date, close_price, volume) VALUES ('600519.SH', '2026-02-24', 1468.02, 25000) ON DUPLICATE KEY UPDATE close_price = VALUES(close_price), volume = VALUES(volume);

执行逻辑

场景数据库行为结果
记录不存在INSERT 新记录数据入库
记录已存在,值相同UPDATE 为相同值数据库状态不变
记录已存在,值不同(数据修正)UPDATE 为新值数据更新为最新

价值

  • 断点续传:5000只股票采集到一半断网,重新运行自动识别存过的更新,没存的插入补全。
  • 数据修正:交易所偶尔修正历史数据,重新跑脚本即可自动更新,无需先删表。

12.2 COALESCE 保护模式

-- 宏观数据更新(COALESCE保护) INSERT INTO trade_macro_indicator (indicator_date, cpi_yoy, ppi_yoy, pmi, m2_yoy, data_source) VALUES ('2026-02-28', NULL, -1.3, 50.8, 8.5, 'akshare') ON DUPLICATE KEY UPDATE ppi_yoy = COALESCE(VALUES(ppi_yoy), ppi_yoy), pmi = COALESCE(VALUES(pmi), pmi), m2_yoy = COALESCE(VALUES(m2_yoy), m2_yoy), data_source = VALUES(data_source);

当 CPI 接口调用失败返回 NULL 时,COALESCE 确保数据库中已有的 2.1 不会被覆盖为 NULL。

12.3 批量写入

# 错误做法:循环逐条写入(5000次数据库连接) for row in rows: cursor.execute(INSERT_SQL, row) # 正确做法:批量写入(1次数据库连接,快10+倍) cursor.executemany(INSERT_SQL, rows)
方式5000条数据耗时数据库连接数
逐条 INSERT~50秒5000次
executemany~3秒1次

13. 性能优化策略

13.1 并行采集策略

数据类型线程数说明
行情数据(QMT)8 线程网络请求为主要瓶颈
财务数据(QMT)1 线程QMT 支持批量传入 stock_list
新闻数据(akshare)8 线程需限速防封 IP
研报数据(akshare)4 线程请求间隔略长

13.2 批量处理

操作批量策略性能提升
财务数据下载50只/批,QMT批量API50x
数据库写入executemany()10x+
最新日期查询一条SQL查全部股票5000x

13.3 缓存策略

  • Redis 热数据缓存:当日行情数据缓存到 Redis,TTL 设为当日收盘后 2 小时
  • 本地文件缓存:QMT 下载的数据自动缓存在本地,重复读取无需联网
  • 行业映射缓存:申万一级行业映射关系变化极低,启动时加载一次缓存全周期使用

13.4 性能基准

操作数据量耗时备注
全量A股东K线增量更新~5000只~5分钟8线程
全量A股财务数据首次下载~5000只~15分钟50只/批
全量A股财务数据增量更新~200只~1分钟跳过已采集
全量A股新闻采集~5000只~50分钟8线程
宏观数据更新7项指标~30秒单线程
财经日历更新~30天~20秒单线程

14. 采集调度策略

14.1 日调度时间表

14.2 调度优先级

优先级任务可延迟失败影响
P0行情数据增量更新下游分析完全不可用
P0换手率计算因子计算缺失
P1宏观数据更新1天宏观环境判断滞后
P1利率数据更新1天流动性判断滞后
P2新闻采集2天情绪因子缺失
P2研报采集7天预期差因子缺失
P3财经日历3天事件预警缺失
P3催化剂发现7天Alpha 机会遗漏

14.3 调度实现

from apscheduler.schedulers.background import BackgroundScheduler scheduler = BackgroundScheduler() # 收盘后行情更新 scheduler.add_job( update_daily_market_data, trigger='cron', hour=15, minute=35, id='market_data_update' ) # 夜间宏观数据 scheduler.add_job( update_macro_indicators, trigger='cron', hour=20, minute=0, id='macro_update' ) # 新闻采集 scheduler.add_job( collect_all_news, trigger='cron', hour=20, minute=20, id='news_collect' )

15. 配置管理

15.1 环境变量 (.env)

# QMT 配置 QMT_HOST=127.0.0.1 QMT_PORT=58000 # 数据库配置 DB_HOST=localhost DB_PORT=1433 DB_NAME=zldata DB_USER=zldata_app DB_PASSWORD=your_secure_password # akshare 配置 AKSHARE_REQUEST_TIMEOUT=30 AKSHARE_REQUEST_DELAY=0.5 AKSHARE_MAX_RETRIES=3 # Redis 配置 REDIS_HOST=localhost REDIS_PORT=6379 REDIS_DB=0 # LLM 配置(催化剂发现) DASHSCOPE_API_KEY=sk-your-key LLM_MODEL=qwen-max # 采集配置 MARKET_DATA_WORKERS=8 FINANCIAL_BATCH_SIZE=50 NEWS_WORKERS=8 REPORT_WORKERS=4

15.2 Pydantic Settings 配置类

from pydantic_settings import BaseSettings class DataCollectionSettings(BaseSettings): """数据采集模块配置""" # QMT qmt_host: str = "127.0.0.1" qmt_port: int = 58000 # 数据库 db_host: str = "localhost" db_port: int = 1433 db_name: str = "zldata" db_user: str = "zldata_app" db_password: str = "" # 采集参数 market_data_workers: int = 8 financial_batch_size: int = 50 news_workers: int = 8 report_workers: int = 4 request_timeout: int = 30 request_delay: float = 0.5 max_retries: int = 3 # LLM dashscope_api_key: str = "" llm_model: str = "qwen-max" class Config: env_file = ".env" env_file_encoding = "utf-8" settings = DataCollectionSettings()

15.3 安全原则

  • 永远不要把数据库密码写在 Python 代码里:使用 .env 文件管理所有敏感配置
  • .env 不入库.gitignore 中必须包含 .env
  • 提供 .env.example:包含所有需要的配置项但不包含真实值,方便新开发者上手

16. 异常处理与监控

16.1 异常分类与处理

异常类型示例处理策略
网络超时QMT 连接超时、akshare HTTP 超时指数退避重试(1s → 2s → 4s → 8s)
数据为空API 返回空 DataFrame切换备用数据源
格式异常日期格式不匹配、字段缺失安全转换 + 默认值 + 记录告警
接口封禁akshare 请求频率过高被限降低线程数、增加请求间隔、暂停10分钟
数据库异常连接断开、死锁重试 3 次后放入重试队列
QMT 不可用miniQMT 未启动切换 akshare + P1 告警

16.2 日志规范

使用 structlog 输出结构化 JSON 日志,包含完整上下文信息:

import structlog logger = structlog.get_logger() # 采集成功日志 logger.info("market_data_downloaded", stock_code="600519.SH", date_range=["2026-02-24", "2026-03-24"], record_count=21, data_source="qmt", duration_ms=1250 ) # 采集失败日志 logger.error("data_source_unavailable", stock_code="000001.SZ", data_source="qmt", error="Connection refused", fallback="switched to akshare" )

16.3 采集监控指标

指标名类型说明告警阈值
dc.job.duration_msHistogram采集任务耗时P99 > 300000ms (5min)
dc.records.writtenCounter成功写入记录数
dc.records.failedCounter写入失败记录数> 0 触发告警
dc.records.skippedCounter跳过记录数(已存在)
dc.source.switch_countCounter数据源切换次数> 10 触发告警
dc.data.freshness_hoursGauge数据新鲜度(小时)> 4 小时触发告警
dc.quality.completenessGauge数据完整率 (%)< 95% 触发告警
dc.quality.anomaly_countCounter异常值数量> 0 触发人工检查

16.4 告警规则

告警条件级别处理方式
行情数据未更新收盘后 1 小时未完成增量更新P0立即通知
双源均不可用QMT + akshare 均返回空/超时P0立即通知 + 自动重试
采集任务失败率 > 5%失败股票数 / 总数 > 5%P11小时内处理
数据源频繁切换1小时内切换 > 10 次P1检查 QMT 服务状态
数据质量下降完整率 < 95% 或异常值增多P2下一工作日排查
新闻/研报采集延迟超过调度时间 2 小时未完成P3下次调度自动补采

17. v1.0 互动展示前端设计

为了向合作方直观展示数据采集模块的建设成果,我们将在 v1.0 版本中,为主平台 zldata.ai-time.net 开发一个 数据采集仪表盘 (Data Collection Dashboard)

17.1 设计目标

  • 可观测性: 直观展示数据管道的健康状况。
  • 可验证性: 允许用户预览已采集的数据,建立对数据质量的信任。
  • 产品化感知: 将”后端脚本”转化为一个”可交互的软件产品”。

17.2 页面布局

17.3 功能组件详情

17.3.1 数据源状态卡片

  • 功能: 实时展示与各主要数据源的连接状态。
  • 数据来源: 调用后端 API /api/v1/dc/status
  • 状态定义:
    • 🟢 在线: 连接测试成功,最近一次采集任务成功。
    • 🟡 降级: 主数据源不可用,已切换至备用数据源。
    • 🔴 离线: 所有数据源均不可用,需人工介入。
  • 交互: 点击卡片可查看详情,如最后错误信息。

17.3.2 最近采集任务日志

  • 功能: 以表格形式展示最近 20 条数据采集任务的执行记录。
  • 数据来源: 调用后端 API /api/v1/dc/jobs
  • 表格字段: 任务类型 (Market/Financial/News) | 开始时间 | 状态 | 受影响行数 | 耗时 (ms) | 详细信息。
  • 状态样式: 成功 (绿色)、失败 (红色)、进行中 (蓝色旋转图标)。

17.3.3 SQL 数据浏览器

这是最核心的互动展示功能,允许用户直接查看数据库中的数据。

  • 功能: 提供一个受限的 SQL 查询控制台,用于探索数据。
  • 安全性设计:
    1. 只读模式: 后端 API (/api/v1/data-viewer/query) 仅解析和执行 SELECT 语句。
    2. 关键词黑名单: 拒绝包含 DROP, DELETE, INSERT, UPDATE, ALTER, CREATE 等关键词的语句。
    3. 表白名单: 仅允许查询 trade_stock_daily, trade_stock_financial, trade_macro_indicator 等数据采集模块产出的表。
    4. 返回限制: 强制添加 LIMIT 子句,最大返回行数不超过 1000 行,防止前端渲染卡死。
    5. 执行超时: 数据库查询设置 5 秒超时。
  • 交互体验:
    • 提供一个带有语法高亮的 SQL 编辑器(可使用 monaco-editorcodemirror)。
    • 提供几个预置的查询模板(如”查询贵州茅台近10日行情”、“查看最新宏观数据”),方便非技术用户使用。
    • 查询结果以分页表格形式展示,支持排序和导出为 CSV。

17.4 前后端接口定义

为支持此仪表盘,后端需提供以下 API 接口。

GET /api/v1/dc/status

响应示例:

{ "qmt": { "status": "online", "last_check": "2026-04-15T08:30:00Z" }, "akshare": { "status": "online", "last_check": "2026-04-15T08:30:00Z" }, "database": { "status": "online", "last_check": "2026-04-15T08:30:00Z" } }

GET /api/v1/dc/jobs

响应示例:

{ "jobs": [ { "id": "job_12345", "type": "market", "start_time": "2026-04-15T15:35:00Z", "end_time": "2026-04-15T15:40:00Z", "status": "success", "affected_rows": 5000, "duration_ms": 300000, "message": "All stocks updated successfully" } ] }

POST /api/v1/data-viewer/query

请求体:

{ "sql": "SELECT stock_code, trade_date, close_price FROM trade_stock_daily WHERE stock_code = '600519.SH' ORDER BY trade_date DESC LIMIT 10;" }

成功响应:

{ "columns": ["stock_code", "trade_date", "close_price"], "rows": [ ["600519.SH", "2026-04-14", 1680.50], ["600519.SH", "2026-04-13", 1675.20] ], "row_count": 2, "execution_time_ms": 45 }

失败响应:

{ "error": "Forbidden query. Only SELECT statements on allowed tables are permitted." }

17.5 开发与部署

  • 前端: 该仪表盘将作为 Next.js 主站 (zldata.ai-time.net) 的一个独立路由 /dashboard/data-collection 进行开发。
  • 后端: FastAPI 服务需要实现上述新增的 API,并部署到 Render Web Service。
  • 部署顺序: 后端 API 先行上线并测试通过,然后前端页面再调用上线后的 API。
  • 技术栈: 前端使用 Next.js 14+ (App Router),后端使用 FastAPI (Python 3.11+)。
  • 安全措施: 所有 API 调用需要进行 CORS 配置,仅允许 zldata.ai-time.net 域名访问。
反馈与评论

来源路径:docs:/docs/backend/data-collection-design

请勿提交账户、密钥、未脱敏交易数据、真实持仓明细或其他敏感信息。

打开反馈服务
Last updated on