02_股票量化数据获取与质量保障体系:从 Tushare 到本地高可用存储
数据是量化交易的基础。本文深入剖析 GPFX 系统如何构建一套完整的股票数据获取、校验、存储、清理体系,涵盖 Tushare API 封装、数据质量校验、复权计算、滚动窗口清理等核心技术,确保数据准确可靠。
一、数据获取架构设计
1.1 整体数据流
Tushare API → 数据校验 → 本地 SQLite → 复权计算 → 质量报告 ↓ ↓ ↓ ↓ ↓ 原始数据 完整性检查 基础存储 衍生数据 问题告警1.2 Tushare API 封装
Tushare 的 pro_api 存在频率限制和偶发异常,需要稳健封装:
classDataService:"""Tushare 数据服务封装"""# 自定义异常类型,区分不同错误场景classTushareRateLimitError(Exception):"""频率超限:需等待后重试"""passclassTushareEmptyResponseError(Exception):"""开市日返回空数据:可能是上游瞬时异常"""passclassTushareIncompleteError(Exception):"""获取完成后仍有交易日全天无数据"""pass@staticmethoddef_retry(func,max_retries:int=3,base_delay:float=1.0):"""指数退避重试策略 - 频率超限:直接抛出,由上层决定等待策略 - 其他异常:指数退避重试,1s → 2s → 4s """forattemptinrange(max_retries):try:returnfunc()exceptExceptionase:msg=str(e)# 频率超限不重试,直接抛出if"频率超限"inmsgor"rate limit"inmsg.lower():raiseTushareRateLimitError(msg)ifattempt<max_retries-1:delay=base_delay*(2**attempt)time.sleep(delay)else:raise设计要点:
- 区分可重试错误(网络超时、空响应)与不可重试错误(频率超限)
- 指数退避避免加重服务端压力
- 自定义异常便于上层精细化处理
二、数据质量保障体系
2.1 三层校验机制
GPFX 实现了从粗到细的三层数据校验:
defverify_daily_data(self,trade_date:str)->DataQualityReport:"""日线数据质量校验"""# 第一层:覆盖度检查 —— 全天是否所有股票都有数据coverage=self._check_coverage(trade_date)# 第二层:完整性检查 —— 关键字段是否缺失completeness=self._check_completeness(trade_date)# 第三层:一致性检查 —— 复权数据与因子是否匹配consistency=self._check_consistency(trade_date)returnDataQualityReport(coverage,completeness,consistency)2.1.1 覆盖度检查
检查每个交易日是否全市场数据完整:
def_check_coverage(self,trade_date:str)->CoverageResult:"""检查指定交易日数据覆盖度"""# 从交易日历获取当日应交易的股票数量expected=self._dao.get_trade_cal_expected_count(trade_date)# 实际获取到的股票数量actual=self._dao.get_daily_count(trade_date,adj_type="none")# 找出缺失的股票missing=self._dao.get_missing_stocks(trade_date)returnCoverageResult(expected=expected,actual=actual,missing_stocks=missing,is_complete=(actual>=expected*0.95)# 允许 5% 停牌)2.1.2 完整性检查
检查关键字段是否为空或异常:
def_check_completeness(self,trade_date:str)->CompletenessResult:"""检查数据完整性"""issues=[]# 检查收盘价是否为 0 或 Nonezero_close=self._dao.query("SELECT COUNT(*) FROM daily_none WHERE trade_date = ? AND (close IS NULL OR close = 0)",(trade_date,))ifzero_close>0:issues.append(f"{zero_close}只股票收盘价异常")# 检查成交量是否为负negative_vol=self._dao.query("SELECT COUNT(*) FROM daily_none WHERE trade_date = ? AND vol < 0",(trade_date,))ifnegative_vol>0:issues.append(f"{negative_vol}只股票成交量为负")returnCompletenessResult(issues=issues,is_valid=len(issues)==0)2.2 复权数据计算与校验
复权是股票数据分析的核心,GPFX 采用本地计算策略:
defcompute_adjusted_data(self,ts_code:str,start_date:str,end_date:str):"""计算前复权/后复权数据 公式: - 前复权价格 = 原始价格 × (当前因子 / 基准日因子) - 后复权价格 = 原始价格 × (当前因子 / 上市首日因子) - 复权成交量 = 原始成交量 × (基准日因子 / 当前因子) """# 获取不复权原始数据raw_df=self._dao.get_daily(ts_code,start_date,end_date,adj_type="none")# 获取复权因子factor_df=self._dao.get_adj_factor(ts_code,start_date,end_date)# 合并数据df=raw_df.merge(factor_df,on=["ts_code","trade_date"],how="left")# 前向填充因子(处理因子缺失)df["adj_factor"]=df["adj_factor"].ffill()# 计算前复权latest_factor=df["adj_factor"].iloc[-1]df["qfq_close"]=df["close"]*(df["adj_factor"]/latest_factor)df["qfq_vol"]=df["vol"]*(latest_factor/df["adj_factor"])# 计算后复权first_factor=df["adj_factor"].iloc[0]df["hfq_close"]=df["close"]*(df["adj_factor"]/first_factor)df["hfq_vol"]=df["vol"]*(first_factor/df["adj_factor"])# 重算涨跌幅(基于复权价)df["qfq_pct_chg"]=df["qfq_close"].pct_change()*100df["hfq_pct_chg"]=df["hfq_close"].pct_change()*100returndf关键细节:
- 价格正向调整(乘以比例),成交量反向调整(除以比例),保持成交额守恒
- 涨跌幅基于复权价重算,避免除权日涨跌幅失真
- 因子前向填充处理偶发缺失
三、退市股处理策略
退市股数据获取是常见痛点,Tushare 对退市股支持有限:
deffetch_delisted_stock_data(self,ts_code:str):"""获取退市股历史数据 策略: 1. 检查是否已标记为"已尝试获取" 2. 调用 Tushare 获取,无论返回多少都标记为已尝试 3. 避免重复调用浪费 API 配额 """# 检查标记ifself._dao.is_delisted_fetched(ts_code):returnFetchResult.SKIPPEDtry:# 尝试获取df=self._api.daily(ts_code=ts_code,start_date="19900101")ifdfisnotNoneandnotdf.empty:self._dao.upsert_daily(ts_code,df,adj_type="none")rows=len(df)else:rows=0# 标记为已尝试(无论成功失败)self._dao.mark_delisted_fetched(ts_code)returnFetchResult(rows=rows,status="fetched")exceptExceptionase:# 也标记为已尝试,避免死循环重试self._dao.mark_delisted_fetched(ts_code)returnFetchResult(rows=0,status=f"failed:{e}")核心思想:退市股数据有限且不会变化,标记机制避免无限重试。
四、滚动窗口数据清理
4.1 清理策略设计
GPFX 采用时间滚动窗口 + 数量上限双策略:
classCleanupService:"""数据清理服务"""def__init__(self,db:Database):self._db=db self._config=get_config()defrun_all(self)->CleanupResult:"""执行全部清理"""result=CleanupResult()# 1. 日线数据:保留最近 N 年(时间窗口)result.daily_none_deleted=self._cleanup_daily_by_years("none")result.daily_qfq_deleted=self._cleanup_daily_by_years("qfq")result.daily_hfq_deleted=self._cleanup_daily_by_years("hfq")# 2. 复权因子/交易日历:同日线窗口result.adj_factor_deleted=self._cleanup_adj_factor()result.trade_cal_deleted=self._cleanup_trade_cal()# 3. 选股/回测结果:保留最近 N 条(数量上限)result.selection_deleted=self._cleanup_selection_by_count()result.backtest_deleted=self._cleanup_backtest_by_count()returnresult4.2 时间窗口锚点设计
清理窗口的锚点必须与获取窗口一致,否则会出现"刚获取就删除"的死循环:
def_resolve_reference_date(self)->str:"""解析清理窗口的参考日期 优先级: 1. 交易日历中最近的开市日(最准确) 2. 日线数据中最近的交易日期 3. 今天(兜底) """# 从交易日历获取最近开市日recent_open=self._dao.get_recent_open_date()ifrecent_open:returnrecent_open# 从日线数据获取recent_daily=self._dao.get_recent_trade_date()ifrecent_daily:returnrecent_daily# 兜底:今天returndatetime.now().strftime("%Y%m%d")为什么用交易日历而非"今天"?
假设今天是 2026-10-06(国庆假期),最近交易日是 2026-09-30:
- 如果用"今天"减 1 年:保留 2025-10-06 之后的数据,会误删 2025-09-30(交易日)
- 如果用最近交易日减 1 年:保留 2025-09-30 之后的数据,边界正确
五、数据预扫描与智能补全
5.1 预扫描机制
获取数据前先扫描缺口,避免重复下载:
defscan_daily_gaps(self,start_date:str,end_date:str)->ScanResult:"""扫描日线数据缺口"""# 获取区间内所有交易日trade_dates=self._dao.get_trade_dates(start_date,end_date)gaps={"none":[],"qfq":[],"hfq":[],}fordateintrade_dates:foradj_typein["none","qfq","hfq"]:count=self._dao.get_daily_count(date,adj_type)expected=self._dao.get_trade_cal_expected_count(date)ifcount<expected*0.95:# 允许 5% 停牌gaps[adj_type].append(date)returnScanResult(gaps=gaps,total_days=len(trade_dates))5.2 智能补全流程
预扫描 → 用户确认 → 后台获取 → 质量校验 → 自动清理 ↓ ↓ ↓ ↓ ↓ 缺口报告 显示进度 批量下载 三层校验 滚动窗口deffetch_missing_data(self,gaps:ScanResult,on_progress=None):"""智能补全缺失数据"""total=sum(len(dates)fordatesingaps.gaps.values())fetched=0foradj_type,datesingaps.gaps.items():fordateindates:# 获取单日全市场数据df=self._fetch_daily_by_date(date,adj_type)ifdfisnotNoneandnotdf.empty:self._dao.upsert_daily_batch(df,adj_type)fetched+=1# 进度回调ifon_progress:on_progress(fetched,total)# 获取完成后校验report=self.verify_daily_data(end_date)ifnotreport.is_complete:raiseTushareIncompleteError(f"获取完成后仍有缺口:{report.missing_dates}")六、性能优化实践
6.1 批量插入
单条插入 5000 只股票需要 30+ 秒,批量插入仅需 1 秒:
defupsert_daily_batch(self,df:pd.DataFrame,adj_type:str):"""批量插入/更新日线数据"""table=f"daily_{adj_type}"# 使用 INSERT OR REPLACE 实现 upsertsql=f""" INSERT OR REPLACE INTO{table}(ts_code, trade_date, open, high, low, close, pre_close, change, pct_chg, vol, amount, ah_vol, ah_amount) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """# 转换为元组列表data=[tuple(row)forrowindf.itertuples(index=False)]# 批量执行withself._db.transaction()asconn:conn.executemany(sql,data)6.2 索引优化
高频查询字段建立索引:
-- 日线表:按股票代码和日期查询最多CREATEINDEXidx_daily_none_codeONdaily_none(ts_code);CREATEINDEXidx_daily_none_dateONdaily_none(trade_date);-- 联合索引:同时按代码和日期范围查询CREATEINDEXidx_daily_none_code_dateONdaily_none(ts_code,trade_date);七、总结
GPFX 的数据获取与质量保障体系核心设计:
| 层级 | 技术方案 | 解决问题 |
|---|---|---|
| API 封装 | 指数退避重试、异常分类 | 网络不稳定、频率限制 |
| 数据校验 | 覆盖度/完整性/一致性三层检查 | 数据缺失、异常值 |
| 复权计算 | 本地因子计算、价量分别调整 | 除权日价格断裂 |
| 退市股 | 标记机制、避免重复 | API 配额浪费 |
| 数据清理 | 时间窗口+数量上限、交易日锚点 | 存储膨胀、窗口不一致 |
| 智能补全 | 预扫描缺口、批量插入 | 重复下载、性能瓶颈 |
这套体系确保了 5000+ 只股票、10+ 年历史数据的准确性和可用性,为上层量化分析提供可靠基础。
关键技术点:Tushare API / SQLite 批量操作 / 复权因子计算 / 数据质量校验 / 滚动窗口清理