量化投资数据仓库构建:从Tushare接口到Parquet存储的工程化实践
1. 项目概述:从零到一构建你的本地股票数据仓库
做量化研究、策略回测,甚至是日常的技术分析,第一步也是最关键的一步,永远是数据。没有高质量、结构化的数据,再精妙的模型也只是空中楼阁。今天,我们不谈那些复杂的算法,就扎扎实实地聊透一个最基础、却让无数新手和老手都踩过坑的环节:如何高效、稳定、自动化地获取股票数据,并把它规规矩矩地保存到本地,形成一个随时可用的数据仓库。
你可能在网上搜到过无数代码片段,用
tushare
、
akshare
或者
baostock
调个接口,
pandas
一存就完事了。但真到实战中,你会发现问题接踵而至:数据源突然失效怎么办?网络波动导致下载中断怎么续传?日积月累的数据如何高效管理和快速读取?不同来源的数据格式如何统一?这篇文章,我将基于我多年管理量化数据集的实战经验,为你拆解一个健壮的股票数据获取与存储系统的核心设计思路和实现细节。这不仅仅是一段脚本,而是一套可扩展的工程化解决方案。
2. 数据源选型与核心接口深度解析
选择数据源是万里长征的第一步,它直接决定了数据的质量、获取成本和后续维护的复杂度。我们主要讨论免费、相对稳定的数据源,并分析其优劣。
2.1 主流免费数据源横向对比
市面上常见的Python数据接口库各有侧重,下表是一个核心对比:
| 数据接口库 | 主要特点 | 数据质量与稳定性 | 获取限制 | 适用场景 |
|---|---|---|---|---|
| Tushare Pro | 数据全面,涵盖股票、基金、期货、宏观经济等,文档专业。 | 较高,有专业团队维护,数据经过清洗。 | 积分制,基础数据免费但需注册获取Token,高频数据需要积分。 | 对数据质量要求较高的量化研究、基本面分析。 |
| AKShare | 接口极其丰富,数据源聚合(来自东方财富、新浪、网易等),开源活跃。 | 取决于原始网站,稳定性一般,可能随网站改版而失效。 | 通常无硬性限制,但需注意爬虫礼仪,避免高频请求。 | 需要非常规数据(如龙虎榜、资金流)、喜欢折腾和贡献的开源爱好者。 |
| Baostock | 专注于A股历史数据,提供除权除息数据,无需Token。 | 稳定,数据来源为官方披露信息,但更新可能有延迟。 | 无明确限制,适合批量获取历史数据。 | 专注于A股历史行情回测,需要干净、准确的复权价格。 |
| Yahoo Finance (yfinance) | 全球市场数据,包括美股、港股、ETF等,接口简洁。 | 对于美股数据质量很好,A股数据为港股ADR,有时有异常。 | 有一定频率限制,大量请求可能被临时阻断。 | 需要美股、全球资产数据的研究者。 |
实操心得 :对于国内A股市场,我建议以 Tushare Pro 作为主力数据源,用 Baostock 作为备份和交叉验证源。Tushare的数据结构规整,社区支持好,虽然需要注册和Token,但免费额度对于个人研究和低频策略完全足够。AKShare可以作为“瑞士军刀”,在需要特定数据时使用,但其接口稳定性需要自己写容错代码来保障。
2.2 接口使用核心细节与避坑指南
以 Tushare Pro 为例,获取日线行情数据看似简单,但细节决定成败。
import tushare as ts
import pandas as pd
# 1. 初始化:Token不要硬编码在代码里!
# 错误示范:pro = ts.pro_api('your_token_here')
# 正确做法:使用环境变量或配置文件
import os
from dotenv import load_dotenv
load_dotenv() # 从 .env 文件加载环境变量
TOKEN = os.getenv('TUSHARE_TOKEN')
if not TOKEN:
raise ValueError("请在 .env 文件中设置 TUSHARE_TOKEN 环境变量")
pro = ts.pro_api(TOKEN)
# 2. 获取数据:理解关键参数
def fetch_daily_data(ts_code, start_date, end_date):
"""
获取复权因子数据,用于计算复权价格
"""
# 先获取复权因子
adj_factor_df = pro.adj_factor(ts_code=ts_code, start_date=start_date, end_date=end_date)
# 获取前复权行情数据
df = pro.daily(ts_code=ts_code, start_date=start_date, end_date=end_date, adj='qfq')
# 关键步骤:合并复权因子,用于验证和自定义复权计算
if not adj_factor_df.empty and not df.empty:
# 确保日期格式一致并合并
adj_factor_df['trade_date'] = pd.to_datetime(adj_factor_df['trade_date']).dt.strftime('%Y%m%d')
df = pd.merge(df, adj_factor_df[['trade_date', 'adj_factor']], on='trade_date', how='left')
return df
# 示例:获取贵州茅台2023年数据
df_maotai = fetch_daily_data('600519.SH', '20230101', '20231231')
print(df_maotai.head())
关键点解析与避坑 :
-
Token管理
:绝对不要将Token直接写在脚本中并上传到Git等公共平台。使用
.env文件配合python-dotenv是行业最佳实践。.env文件应加入.gitignore。 -
复权方式
:
adj='qfq'代表前复权,这是回测最常用的格式,价格与当前价格可比。adj='hfq'是后复权,总市值保持一致。务必明确你的策略需要哪种数据。 -
日期格式
:Tushare接口的日期参数格式通常是
'YYYYMMDD'的字符串,而返回的DataFrame里的trade_date列也是该格式。进行时间序列分析时,需要将其转换为datetime类型:df['trade_date'] = pd.to_datetime(df['trade_date'])。 - 网络超时与重试 :批量获取时网络问题不可避免。必须封装带有重试机制的请求函数。
import time
from tenacity import retry, stop_after_attempt, wait_exponential
@retry(stop=stop_after_attempt(5), wait=wait_exponential(multiplier=1, min=4, max=10))
def robust_fetch_data(api_func, **kwargs):
"""
带重试机制的数据获取函数
"""
try:
df = api_func(**kwargs)
# 如果返回数据为空,可能是代码错误或日期错误,不重试
if df is None or df.empty:
print(f"警告: 查询参数 {kwargs} 返回空数据。")
return df
return df
except Exception as e:
print(f"请求失败: {e}, 参数: {kwargs}. 准备重试...")
raise e # 触发tenacity重试
3. 数据存储方案设计与工程化实践
获取到数据只是开始,如何存储决定了未来使用的效率。我们追求的是: 读写快、易管理、可追溯、能扩展 。
3.1 存储格式选型:从CSV到数据库
| 存储格式 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|
| CSV | 人类可读,通用性强,任何工具都能打开。 | 读写速度慢(尤其大文件),无数据类型校验,修改效率低。 | 小型数据集,临时交换数据,需要给人看的报告。 |
| Parquet | 列式存储,压缩率高,读写速度极快(特别适合pandas),支持复杂数据类型。 | 文件格式二进制,需要特定库(如pyarrow)读取。 | 强烈推荐 作为本地主力存储格式,用于历史数据归档。 |
| Feather | 读写速度最快,设计用于pandas DataFrame的快速序列化。 | 格式不如Parquet通用,压缩率相对较低。 | 需要极速读写的中间临时数据。 |
| SQLite | 单个文件数据库,支持SQL查询,具备事务、索引等功能。 | 并发写入性能有瓶颈,不适合超高频写入。 | 中小型项目,需要复杂查询和关系管理的数据。 |
| MySQL/PostgreSQL | 功能完整的关系型数据库,强大的查询和管理能力。 | 需要单独部署和维护,架构复杂。 | 团队协作,数据量极大,需要严格事务和权限管理的生产环境。 |
对于个人或小型量化项目,我的推荐组合是: 使用 Parquet 文件按股票代码分目录存储历史数据,使用 SQLite 存储元信息(如股票列表、更新状态) 。
3.2 本地文件系统组织结构设计
一个清晰的文件结构是高效管理的基础。切忌所有数据扔进一个文件夹。
stock_data/
├── meta/ # 元数据
│ ├── stock_basic.parquet # 股票基本信息表
│ └── update_log.db # SQLite,记录各股票数据更新日期
├── daily/ # 日线数据
│ ├── SH/ # 上海交易所
│ │ ├── 600519.parquet # 贵州茅台
│ │ └── 600036.parquet # 招商银行
│ └── SZ/ # 深圳交易所
│ ├── 000001.parquet # 平安银行
│ └── 300750.parquet # 宁德时代
├── adj_factor/ # 复权因子单独存储(可选)
│ ├── SH/
│ └── SZ/
└── scripts/ # 数据维护脚本
├── downloader.py
├── updater.py
└── validator.py
为什么按代码和交易所分目录?
- 并行化处理 :下载或更新时,可以按股票代码并行,互不干扰。
- 快速定位 :根据代码直接定位文件路径,无需遍历或查询数据库。
- 备份灵活 :可以轻松备份或同步单个股票或整个交易所的数据。
3.3 使用Parquet进行高效存储的代码实现
import pandas as pd
import pyarrow as pa
import pyarrow.parquet as pq
import os
from pathlib import Path
class ParquetDataStore:
def __init__(self, base_path='./stock_data'):
self.base_path = Path(base_path)
self.daily_path = self.base_path / 'daily'
# 初始化目录
self.daily_path.mkdir(parents=True, exist_ok=True)
(self.daily_path / 'SH').mkdir(exist_ok=True)
(self.daily_path / 'SZ').mkdir(exist_ok=True)
def _get_file_path(self, ts_code, data_type='daily'):
"""根据股票代码生成文件路径"""
# 示例代码: 600519.SH -> SH/600519.parquet
code, exchange = ts_code.split('.')
if exchange not in ['SH', 'SZ']:
raise ValueError(f"不支持的交易所代码: {exchange}")
if data_type == 'daily':
return self.daily_path / exchange / f"{code}.parquet"
# 可以扩展其他数据类型路径
else:
return self.base_path / data_type / exchange / f"{code}.parquet"
def save_daily_data(self, ts_code, df):
"""保存单只股票的日线数据到Parquet"""
if df.empty:
print(f"{ts_code}: 数据为空,跳过保存。")
return
file_path = self._get_file_path(ts_code, 'daily')
# 关键操作:将trade_date设为索引,并确保排序
df = df.copy()
df['trade_date'] = pd.to_datetime(df['trade_date'])
df.set_index('trade_date', inplace=True)
df.sort_index(inplace=True) # 按日期升序排列
# 使用pyarrow写入,指定压缩方式以节省空间
table = pa.Table.from_pandas(df, preserve_index=True)
pq.write_table(table, file_path, compression='snappy') # snappy压缩速度快
print(f"数据已保存至: {file_path}")
def load_daily_data(self, ts_code, start_date=None, end_date=None):
"""从Parquet加载单只股票日线数据,支持日期切片"""
file_path = self._get_file_path(ts_code, 'daily')
if not file_path.exists():
print(f"文件不存在: {file_path}")
return pd.DataFrame()
# 高效读取:可以只读取需要的列和行
table = pq.read_table(file_path)
df = table.to_pandas()
# 日期筛选
if start_date:
start_date = pd.Timestamp(start_date)
df = df[df.index >= start_date]
if end_date:
end_date = pd.Timestamp(end_date)
df = df[df.index <= end_date]
return df
# 使用示例
store = ParquetDataStore()
# 假设 df_maotai 是之前获取的DataFrame
store.save_daily_data('600519.SH', df_maotai)
# 加载2023年3月的数据
loaded_data = store.load_daily_data('600519.SH', start_date='2023-03-01', end_date='2023-03-31')
注意事项 :
-
索引设置
:将
trade_date设为DataFrame的索引是时间序列数据分析的标准操作,能极大提升按日期查询和合并的效率。 - 排序 :保存前务必按日期排序,保证数据一致性。
-
压缩
:Parquet支持多种压缩算法(如
snappy,gzip)。snappy压缩和解压速度极快,占用CPU少,是平衡速度和空间的好选择。gzip压缩率更高,但更耗CPU。 - 分区 :对于超大数据集,Parquet支持按列(如年份、月份)进行分区存储,能进一步提升查询性能。但对于单股票数据,通常不需要。
4. 自动化更新与增量同步策略
数据不是一次性下载就完事的,需要持续更新。一个健壮的更新机制需要解决: 识别缺失数据、处理网络异常、避免重复下载、记录更新状态 。
4.1 增量更新逻辑设计
核心思想:本地已有什么数据,就只下载缺失的数据。
import sqlite3
from datetime import datetime, timedelta
class DataUpdater:
def __init__(self, store: ParquetDataStore, db_path='./stock_data/meta/update_log.db'):
self.store = store
self.db_path = Path(db_path)
self.db_path.parent.mkdir(parents=True, exist_ok=True)
self._init_db()
def _init_db(self):
"""初始化SQLite数据库,创建更新记录表"""
conn = sqlite3.connect(self.db_path)
cursor = conn.cursor()
cursor.execute('''
CREATE TABLE IF NOT EXISTS update_log (
ts_code TEXT PRIMARY KEY,
last_update_date TEXT, -- 最后更新到的交易日
last_attempt_date TEXT, -- 最后尝试更新日期
status TEXT -- 状态: success, failed, pending
)
''')
conn.commit()
conn.close()
def get_local_latest_date(self, ts_code):
"""获取本地该股票最新的数据日期"""
try:
df = self.store.load_daily_data(ts_code)
if df.empty:
return None
# 索引是trade_date,取最大值
return df.index.max().strftime('%Y%m%d')
except FileNotFoundError:
return None
def calculate_update_range(self, ts_code):
"""计算需要更新的日期范围"""
local_latest = self.get_local_latest_date(ts_code)
# 如果本地没有数据,则从头开始下载(例如从上市日期或一年前开始)
if local_latest is None:
# 这里需要调用接口获取股票上市日期,简化处理,假设从一年前开始
start_date = (datetime.now() - timedelta(days=365)).strftime('%Y%m%d')
else:
# 本地最新日期的下一天作为开始
latest_dt = datetime.strptime(local_latest, '%Y%m%d')
start_date = (latest_dt + timedelta(days=1)).strftime('%Y%m%d')
end_date = datetime.now().strftime('%Y%m%d')
# 还需要判断start_date是否晚于end_date(即数据已是最新)
if start_date > end_date:
return None, None # 无需更新
return start_date, end_date
def update_single_stock(self, ts_code):
"""更新单只股票数据"""
start_date, end_date = self.calculate_update_range(ts_code)
if start_date is None:
print(f"{ts_code}: 本地数据已是最新,无需更新。")
return
print(f"{ts_code}: 准备下载 {start_date} 至 {end_date} 的数据。")
try:
# 使用前面封装好的带重试的获取函数
new_df = robust_fetch_data(pro.daily, ts_code=ts_code, start_date=start_date, end_date=end_date, adj='qfq')
if new_df is not None and not new_df.empty:
# 加载本地已有数据,合并新数据,去重后保存
local_df = self.store.load_daily_data(ts_code)
combined_df = pd.concat([local_df, new_df.set_index('trade_date')]).sort_index()
# 基于索引去重,保留最后出现的(即新数据)
combined_df = combined_df[~combined_df.index.duplicated(keep='last')]
self.store.save_daily_data(ts_code, combined_df.reset_index())
# 更新日志
self._log_update(ts_code, end_date, 'success')
print(f"{ts_code}: 成功更新至 {end_date}。")
else:
print(f"{ts_code}: 未获取到新数据。")
self._log_update(ts_code, datetime.now().strftime('%Y%m%d'), 'failed')
except Exception as e:
print(f"{ts_code}: 更新失败,错误: {e}")
self._log_update(ts_code, datetime.now().strftime('%Y%m%d'), 'failed')
def _log_update(self, ts_code, update_date, status):
"""记录更新状态到数据库"""
conn = sqlite3.connect(self.db_path)
cursor = conn.cursor()
attempt_date = datetime.now().strftime('%Y%m%d %H:%M:%S')
cursor.execute('''
INSERT OR REPLACE INTO update_log (ts_code, last_update_date, last_attempt_date, status)
VALUES (?, ?, ?, ?)
''', (ts_code, update_date if status == 'success' else None, attempt_date, status))
conn.commit()
conn.close()
4.2 批量更新与任务调度
有了单股票更新能力,批量更新就简单了。关键在于 控制并发和频率 ,避免给数据源服务器造成压力。
import concurrent.futures
import time
def batch_update_stocks(ts_code_list, max_workers=5):
"""
批量更新股票列表
max_workers: 控制并发线程数,建议不要超过5,避免被封IP
"""
store = ParquetDataStore()
updater = DataUpdater(store)
def update_task(ts_code):
updater.update_single_stock(ts_code)
# 礼貌性延迟,模拟人工操作
time.sleep(0.5)
return ts_code
with concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) as executor:
future_to_code = {executor.submit(update_task, code): code for code in ts_code_list}
for future in concurrent.futures.as_completed(future_to_code):
code = future_to_code[future]
try:
future.result()
except Exception as exc:
print(f'{code} generated an exception: {exc}')
# 示例:更新一批股票
stock_list = ['600519.SH', '000001.SZ', '300750.SZ', '000858.SZ']
batch_update_stocks(stock_list, max_workers=3)
自动化调度 : 对于每日更新,可以结合系统定时任务(如Linux的cron,Windows的任务计划程序)或Python的APScheduler库。
from apscheduler.schedulers.blocking import BlockingScheduler
def scheduled_daily_update():
print(f"开始每日数据更新任务,时间:{datetime.now()}")
# 1. 获取需要更新的股票列表(例如,所有已存储的股票)
# 2. 调用 batch_update_stocks
print("每日更新任务完成。")
if __name__ == '__main__':
scheduler = BlockingScheduler()
# 每个交易日收盘后(例如下午6点)运行
scheduler.add_job(scheduled_daily_update, 'cron', hour=18, minute=0, day_of_week='mon-fri')
print('数据更新调度器已启动,按 Ctrl+C 退出。')
try:
scheduler.start()
except (KeyboardInterrupt, SystemExit):
pass
5. 数据质量校验与常见问题排查
数据错了,一切分析都是白费。建立数据校验机制至关重要。
5.1 常见数据质量问题清单
- 缺失值 :特别是停牌日,数据接口可能返回空或NaN。需要区分是“当日无交易”还是“数据获取失败”。
- 价格异常 :收盘价、开盘价超出涨跌停范围(A股普通股票±10%,ST股±5%),或者出现负值、极大值。
- 复权不一致 :不同数据源对同一只股票的复权价格可能有细微差异。
- 日期不连续 :数据中缺失了交易日(可能是因为接口限制、网络问题)。
- 成交量/成交额为0 :非停牌日出现零成交,可能是数据错误。
5.2 自动化校验脚本实现
def validate_stock_data(df, ts_code):
"""
对单只股票的DataFrame进行基础校验
返回一个包含问题的字典列表
"""
issues = []
if df.empty:
issues.append({'code': ts_code, 'type': 'EMPTY', 'desc': '数据为空'})
return issues
# 1. 检查日期连续性
df_sorted = df.sort_index()
date_diff = df_sorted.index.to_series().diff().dt.days
# 正常情况下,差值应该是1(连续交易日)或更大(间隔了非交易日或节假日)
# 但差值大于3天可能就有问题,需要结合日历判断,这里简化处理
gap_dates = df_sorted.index[date_diff > 3]
if not gap_dates.empty:
for d in gap_dates[:3]: # 只报告前三个缺口
issues.append({'code': ts_code, 'type': 'DATE_GAP', 'date': d.strftime('%Y%m%d'), 'desc': f'日期不连续,与前一日间隔{date_diff.loc[d]}天'})
# 2. 检查价格合理性(简单逻辑)
# 假设股价在1元到10000元之间
price_cols = ['open', 'high', 'low', 'close']
for col in price_cols:
if col in df.columns:
invalid_prices = df[(df[col] <= 0) | (df[col] > 10000)]
if not invalid_prices.empty:
for idx, row in invalid_prices.iterrows():
issues.append({'code': ts_code, 'type': 'PRICE_ABNORMAL', 'date': idx.strftime('%Y%m%d'), 'field': col, 'value': row[col], 'desc': '价格超出合理范围'})
# 3. 检查涨跌幅是否异常(基于前复权价格)
if 'close' in df.columns and 'pre_close' in df.columns:
df['pct_chg_calc'] = (df['close'] - df['pre_close']) / df['pre_close'] * 100
# 与接口提供的pct_chg对比,允许微小浮点误差
if 'pct_chg' in df.columns:
mismatch = df[abs(df['pct_chg_calc'] - df['pct_chg']) > 0.015] # 允许0.015%的误差
if not mismatch.empty:
for idx, row in mismatch.iterrows():
issues.append({'code': ts_code, 'type': 'PCT_CHG_MISMATCH', 'date': idx.strftime('%Y%m%d'), 'desc': f'计算涨跌幅{row[\"pct_chg_calc\"]:.2f}%与接口提供{row[\"pct_chg\"]:.2f}%不符'})
return issues
# 遍历数据目录,校验所有股票
def validate_all_data(store_path='./stock_data/daily'):
all_issues = []
base_path = Path(store_path)
for exchange_dir in base_path.iterdir():
if exchange_dir.is_dir():
for parquet_file in exchange_dir.glob('*.parquet'):
ts_code = f"{parquet_file.stem}.{exchange_dir.name}"
df = pd.read_parquet(parquet_file)
issues = validate_stock_data(df, ts_code)
all_issues.extend(issues)
return all_issues
5.3 网络与接口问题排查实录
问题1:
ts.pro_api().daily()
返回
None
或空
DataFrame
。
- 可能原因1:Token无效或过期。 检查Token是否正确,并在Tushare官网查看积分和调用权限。
-
可能原因2:股票代码格式错误。
必须是
'代码.交易所'格式,如'600519.SH'。SH代表上海,SZ代表深圳。 - 可能原因3:日期范围内无数据。 股票可能尚未上市、已退市,或请求的是非交易日。
-
排查步骤
:
- 打印请求参数确认无误。
-
尝试一个肯定有数据的股票和日期(如
'000001.SZ'和最近一个交易日)。 - 直接访问Tushare官网的测试工具,用相同参数测试。
问题2:批量下载时频繁被中断或封IP。
-
解决方案
:
-
降低并发度
:将
max_workers从 10 降到 3 或更低。 -
增加延迟
:在每个请求之间加入随机延时
time.sleep(random.uniform(1, 3))。 - 使用代理IP池 :对于海量数据抓取,这是终极方案,但维护成本高。对于免费接口,更建议礼貌、低频地调用。
-
降低并发度
:将
问题3:保存的Parquet文件用
pd.read_parquet()
读取时列名或数据类型错误。
- 可能原因 :不同批次的数据字段顺序或类型不一致,导致合并后保存出错。
-
解决方案
:在保存前,强制统一
DataFrame的列顺序和数据类型。
def standardize_dataframe(df, expected_columns):
"""标准化DataFrame的列顺序和类型"""
# 确保包含所有预期列,缺失的用NaN填充
for col in expected_columns:
if col not in df.columns:
df[col] = pd.NA
# 按预期顺序排列列
df = df[expected_columns]
# 可以在这里添加类型转换,例如确保数值列是float
numeric_cols = ['open', 'high', 'low', 'close', 'pre_close', 'change', 'pct_chg', 'vol', 'amount']
for col in numeric_cols:
if col in df.columns:
df[col] = pd.to_numeric(df[col], errors='coerce')
return df
构建本地股票数据仓库是一个系统工程,它远不止于运行几行下载代码。从数据源的谨慎选型、接口的稳健调用,到存储格式的精心设计、目录结构的清晰规划,再到增量更新、质量校验和异常处理,每一个环节都需要考虑周全。这套方法论和代码框架,是我经过多个项目迭代后总结出的相对稳定的实践。它可能不是最完美的,但足够让你避开我当年踩过的大多数坑,为你后续的量化分析和策略研究打下一个可靠的数据地基。记住,干净、完整、易于获取的数据,是任何数据分析工作成功的一半。
更多推荐

所有评论(0)