ZLData 数据采集模块 — 详细设计文档
| 属性 | 内容 |
|---|---|
| 文档编号 | ZLData-DC-DD-001 |
| 配套文档 | ZLData_PRD_v2.md v2.1 |
| 日期 | 2026-04-03 |
| 状态 | Draft — 待评审 |
目录
- 模块概述
- 数据源设计
- 双源冗余架构
- 行情数据采集
- 财务数据采集
- 宏观经济数据采集
- 利率数据采集
- 新闻事件采集与情感分析
- 机构研报采集
- 财经日历与催化剂发现
- 数据标准化与质量保障
- 数据入库设计
- 性能优化策略
- 采集调度策略
- 配置管理
- 异常处理与监控
1. 模块概述
1.1 定位与目标
数据采集模块(Data Collection)是 ZLData AI 量化交易平台的 数据基础设施层,承担着从多源异构数据源获取行情、财务、宏观、新闻和研报数据,并经过清洗、标准化、去重后写入中心数据库的核心职责。该模块是整个平台的数据基石——量化分析、策略回测、风险评估和智能体决策的上游输入完全依赖本模块所提供的数据质量与时效性。
设计目标包括:
- 数据完整性:覆盖技术分析(行情)、基本面分析(财务)、宏观环境(指标)和另类数据(新闻/研报/事件)四类数据,构建完整的 AI 特征库。
- 高可用性:通过双数据源冗余架构,确保单一数据源故障时数据采集不中断,数据可用性达到 99.99%。
- 时效性保障:日内数据延迟 < 5 秒,日频数据在收盘后 30 分钟内完成更新。
- 幂等与可恢复:任何采集任务可安全重复执行,断网或崩溃后从断点续传,无需删表重来。
- 生产级质量:包含完整的异常处理、数据质量检查、限速策略和监控告警机制。
1.2 数据分类体系
1.3 核心数据表清单
| 序号 | 表名 | 说明 | 数据量级 | 更新频率 |
|---|---|---|---|---|
| 1 | trade_stock_daily | 日K线行情 | ~5000只 × 3000日 ≈ 1500万行 | 日频 |
| 2 | trade_stock_financial | 季度财务指标 | ~5000只 × 40期 ≈ 20万行 | 季频 |
| 3 | trade_macro_indicator | 月度宏观指标 | ~120行(10年) | 月频 |
| 4 | trade_rate_daily | 日频利率指标 | ~2500行(10年) | 日频 |
| 5 | trade_stock_news | 个股新闻事件 | ~50万行(年增量) | 小时级 |
| 6 | trade_report_consensus | 机构研报一致性 | ~10万行 | 周频 |
| 7 | trade_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() | 新浪财经 | 利润表/资产负债表/现金流量表 |
| 宏观CPI | ak.macro_china_cpi_monthly() | 国家统计局 | 月度同比/环比/累计 |
| 宏观PPI | ak.macro_china_ppi_monthly() | 国家统计局 | 月度同比/环比 |
| 宏观PMI | ak.macro_china_pmi() | 国家统计局 | 制造业/非制造业PMI |
| 宏观M2 | ak.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 的数据差异处理:
| 维度 | QMT | akshare | 统一策略 |
|---|---|---|---|
| 股票代码格式 | 600519.SH | 600519 | 内部统一为 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 None4.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 张财务报表数据:
| 报表 | 表名 | 关键字段 |
|---|---|---|
| 资产负债表 | Balance | tot_assets, tot_liab, total_equity, total_current_assets, total_current_liability, inventories, cash_equivalents |
| 利润表 | Income | revenue, net_profit_incl_min_int_inc, operating_revenue, oper_profit, cost_of_goods_sold |
| 现金流量表 | CashFlow | net_cash_flows_oper_act, net_cash_flows_inv_act, net_cash_flows_fnc_act |
| 每股指标表 | PershareIndex | s_fa_eps_basic, s_fa_bps, s_fa_ocfps, s_fa_undistributedps, du_return_on_equity, sales_gross_profit |
| 股本表 | Capital | totalShares, 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 | 净利润 / 净资产 × 100 | Income + Balance | 先取 PershareIndex,再计算 |
| ROA | 净利润 / 总资产 × 100 | Income + Balance | 始终计算 |
| 毛利率 | (营收 - 营业成本) / 营收 × 100 | Income | 先取 PershareIndex,再计算 |
| 净利率 | 净利润 / 营业收入 × 100 | Income | 始终计算 |
| 资产负债率 | 总负债 / 总资产 × 100 | Balance | 始终计算 |
| 流动比率 | 流动资产 / 流动负债 | Balance | 始终计算 |
| 速动比率 | (流动资产 - 存货) / 流动负债 | Balance | 始终计算 |
| 总资产周转率 | 营业收入 / 总资产 | Income + Balance | 始终计算 |
| OCF/营收 | 经营现金流 / 营业收入 × 100 | CashFlow + 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() | 月频 | ★★★★★ |
| PMI | ak.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月、202601、2026-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 日期标准化
不同数据源返回不同的日期格式,统一处理策略:
| 原始格式 | 来源 | 标准化结果 |
|---|---|---|
20260331 | QMT (m_timetag) | 2026-03-31 (DATE) |
1709251200000 | QMT (毫秒时间戳) | 2024-03-01 (DATE) |
2026年03月 | akshare 宏观 | 2026-03-31 (月末日期) |
202603 | akshare 宏观 | 2026-03-31 (月末日期) |
2026-03-31 | akshare 行情 | 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 None11.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批量API | 50x |
| 数据库写入 | 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=415.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_ms | Histogram | 采集任务耗时 | P99 > 300000ms (5min) |
dc.records.written | Counter | 成功写入记录数 | — |
dc.records.failed | Counter | 写入失败记录数 | > 0 触发告警 |
dc.records.skipped | Counter | 跳过记录数(已存在) | — |
dc.source.switch_count | Counter | 数据源切换次数 | > 10 触发告警 |
dc.data.freshness_hours | Gauge | 数据新鲜度(小时) | > 4 小时触发告警 |
dc.quality.completeness | Gauge | 数据完整率 (%) | < 95% 触发告警 |
dc.quality.anomaly_count | Counter | 异常值数量 | > 0 触发人工检查 |
16.4 告警规则
| 告警 | 条件 | 级别 | 处理方式 |
|---|---|---|---|
| 行情数据未更新 | 收盘后 1 小时未完成增量更新 | P0 | 立即通知 |
| 双源均不可用 | QMT + akshare 均返回空/超时 | P0 | 立即通知 + 自动重试 |
| 采集任务失败率 > 5% | 失败股票数 / 总数 > 5% | P1 | 1小时内处理 |
| 数据源频繁切换 | 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 查询控制台,用于探索数据。
- 安全性设计:
- 只读模式: 后端 API (
/api/v1/data-viewer/query) 仅解析和执行SELECT语句。 - 关键词黑名单: 拒绝包含
DROP,DELETE,INSERT,UPDATE,ALTER,CREATE等关键词的语句。 - 表白名单: 仅允许查询
trade_stock_daily,trade_stock_financial,trade_macro_indicator等数据采集模块产出的表。 - 返回限制: 强制添加
LIMIT子句,最大返回行数不超过 1000 行,防止前端渲染卡死。 - 执行超时: 数据库查询设置 5 秒超时。
- 只读模式: 后端 API (
- 交互体验:
- 提供一个带有语法高亮的 SQL 编辑器(可使用
monaco-editor或codemirror)。 - 提供几个预置的查询模板(如”查询贵州茅台近10日行情”、“查看最新宏观数据”),方便非技术用户使用。
- 查询结果以分页表格形式展示,支持排序和导出为 CSV。
- 提供一个带有语法高亮的 SQL 编辑器(可使用
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请勿提交账户、密钥、未脱敏交易数据、真实持仓明细或其他敏感信息。