共计 3157 个字符,预计需要花费 8 分钟才能阅读完成。
背景痛点
在 A 股量化研究中,数据采集是策略开发的基石,但实际操作中常遇到以下典型问题:

- 交易所限频:交易所接口对访问频率有严格限制,如沪深交易所的分钟级访问上限,容易触发封禁
- 历史数据缺失:部分免费数据源存在字段不全或历史数据断层问题(如 ST 股票特殊处理状态标记)
- 接口不稳定:网络波动或数据源维护导致接口暂时不可用
- 数据一致性:不同数据源对同一标的的复权处理方式可能不同
对比常见数据源:
- Tushare Pro:
- 优势:数据质量高,包含财务因子等特色数据
-
不足:积分制度限制免费用户调用频次
-
Baostock:
- 优势:无需认证,提供完整的除权除息数据
-
不足:不支持分钟级实时行情
-
AKShare:
- 优势:完全免费、接口丰富(含港股 / 期货)、社区活跃
- 不足:缺乏官方文档,部分接口稳定性依赖第三方网站
技术方案
AKShare 模块化设计
AKShare 采用 ” 数据域 ” 划分模式,主要模块包括:
stock_zh_a:A 股基础行情stock_fundamental:财务报表数据stock_zh_index:指数相关数据
每个模块独立维护数据源适配器,通过统一异常处理层保证模块间隔离。
多进程采集架构
[主进程] → [任务队列] → [Worker 进程 1] → [数据源 A]
→ [Worker 进程 2] → [数据源 B]
→ [Monitor 进程] → [异常报警]
关键组件:
- 任务调度器:将股票代码按行业分组,均衡分配到不同 Worker
- 进程池:控制并发数量(建议为 CPU 核心数×2)
- 结果聚合器:合并各 Worker 返回的 DataFrame
令牌桶限频算法
from threading import Lock
import time
class TokenBucket:
def __init__(self, capacity, fill_rate):
self.capacity = capacity # 桶容量
self._tokens = capacity
self.fill_rate = fill_rate # 令牌 / 秒
self.timestamp = time.time()
self.lock = Lock()
def consume(self):
with self.lock:
now = time.time()
delta = self.fill_rate * (now - self.timestamp)
self._tokens = min(self.capacity, self._tokens + delta)
self.timestamp = now
if self._tokens >= 1:
self._tokens -= 1
return True
return False
代码实现
异步获取股票列表
import akshare as ak
from concurrent.futures import ThreadPoolExecutor
def get_all_stocks():
"""获取沪深全量股票列表(含退市)"""
stock_zh = ak.stock_zh_a_spot()
stock_sh = ak.stock_sh_a_spot()
stock_sz = ak.stock_sz_a_spot()
return pd.concat([stock_zh, stock_sh, stock_sz]).drop_duplicates()
带重试的日 K 线下载
import random
from tenacity import retry, stop_after_attempt, wait_exponential
@retry(stop=stop_after_attempt(3),
wait=wait_exponential(multiplier=1, min=4, max=10))
def get_daily(symbol, start_date, end_date):
try:
# 随机选择数据源降低封禁风险
if random.random() > 0.5:
return ak.stock_zh_a_daily(symbol=symbol, start_date=start_date, end_date=end_date)
else:
return ak.stock_zh_a_hist(symbol=symbol, period="daily", start_date=start_date, end_date=end_date)
except Exception as e:
print(f"{symbol} 数据获取失败: {str(e)}")
raise
SQLite 本地缓存
import sqlite3
from datetime import datetime
def init_db():
conn = sqlite3.connect('quant.db')
c = conn.cursor()
c.execute('''CREATE TABLE IF NOT EXISTS daily
(symbol text, date text, open real, high real,
low real, close real, volume real, PRIMARY KEY (symbol, date))''')
conn.commit()
return conn
def save_to_cache(conn, df, symbol):
df['symbol'] = symbol
df.to_sql('daily', conn, if_exists='append', index=False)
性能优化
并发模式对比
测试环境:4 核 CPU/8GB 内存,采集 2023 年全 A 股日线数据
| 模式 | 耗时 | CPU 利用率 |
|---|---|---|
| 单线程 | 82min | 15% |
| 4 进程 | 19min | 92% |
| 8 进程 + 缓存 | 14min | 98% |
内存监控方案
import tracemalloc
tracemalloc.start()
# ... 执行数据采集...
snapshot = tracemalloc.take_snapshot()
top_stats = snapshot.statistics('lineno')
for stat in top_stats[:5]:
print(stat)
网络异常测试
使用 unittest.mock 模拟网络抖动:
from unittest.mock import patch
import pytest
@patch('akshare.stock_zh_a_daily')
def test_network_failure(mock_api):
mock_api.side_effect = [TimeoutError, None]
result = get_daily("000001", "20230101", "20231231")
assert not result.empty
避坑指南
IP 封禁预防
- 使用代理 IP 池轮询(推荐 Luminati 等商业服务)
- 设置随机延迟:
time.sleep(random.uniform(0.5, 2)) - 避免在交易所结算时间段(15:30-16:30)密集访问
除权除息处理
正确步骤:
- 先获取原始行情数据
- 单独下载
stock_zh_a_cdr除权除息表 - 使用
pandas.merge按除权日对齐
回测数据一致性
关键检查点:
- 验证复权因子计算逻辑(前复权 / 后复权)
- 检查停牌日期是否留有无效数据
- 对比不同数据源的成交量单位(手 / 股)
延伸思考
- 实时行情对接:通过 WebSocket 连接交易所 Level2 行情(需券商 API 权限)
- 数据质量监控:建立字段完整性、时间连续性等检测规则
- 分布式存储:使用 Dask+Parquet 处理超大规模历史数据
实践心得
经过三个月的生产环境运行,这套方案稳定支撑了日均 20 万次的 API 调用。最大的收获是认识到:量化数据系统的可靠性不是单纯的技术问题,而是需要将工程实践(重试机制)、金融知识(除权处理)和运维经验(IP 规避)结合的综合性解决方案。建议开发者在初期就建立完善的数据校验日志,这将为后续的策略回测节省大量调试时间。
正文完
