批量数据获取与并发控制
前一篇完成了环境搭建和单个股票数据获取。这篇解决实际问题:如何高效获取全市场 5000+ 支股票的数据?
问题分析
单线程获取的瓶颈
# 单线程获取
for ts_code in stock_list:
df = client.get_daily_data(ts_code, \’20230101\’, \’20231231\’)
time.sleep(0.1) # 频率控制
# 问题:
# – 5000 支股票 × 0.1 秒 = 500 秒 ≈ 8.3 分钟
# – 如果获取 5 年数据,需要 40+ 分钟
# – 网络异常会导致整个流程中断
需要解决的问题
解决方案:并发控制
方案对比
| 多线程 | 简单,IO 友好 | GIL 限制 | IO 密集型 ✓ |
| 多进程 | 突破 GIL | 内存占用大 | CPU 密集型 |
| 异步 IO | 性能最好 | 学习成本高 | 高并发场景 |
| 进程池 | 简单易用 | 灵活性一般 | 中等规模 ✓ |
选择:进程池(ProcessPoolExecutor)+ 线程池(ThreadPoolExecutor)组合
实现方案
方案 1:线程池(推荐入门)
创建 src/batch_fetch_thread.py:
\”\”\”批量获取日线数据 – 线程池版本\”\”\”
from data_client import TushareClient
from concurrent.futures import ThreadPoolExecutor, as_completed
from pathlib import Path
import time
from typing import List, Tuple
class BatchFetcher:
\”\”\”批量数据获取器\”\”\”
def __init__(self, token: str = None, max_workers: int = 5):
\”\”\”
初始化获取器
Args:
token: Tushare token
max_workers: 最大并发数(默认 5,避免触发限流)
\”\”\”
self.client = TushareClient(token)
self.max_workers = max_workers
self.output_dir = Path(\’data/csv/daily\’)
self.output_dir.mkdir(parents=True, exist_ok=True)
def fetch_single_stock(
self,
args: Tuple[str, str, str]
) –> Tuple[str, bool, str]:
\”\”\”
获取单支股票数据(用于并发执行)
Args:
args: (ts_code, start_date, end_date)
Returns:
(ts_code, success, message)
\”\”\”
ts_code, start_date, end_date = args
try:
# 检查是否已存在
filename = f\”{
ts_code.replace(\’.\’, \’_\’)}.csv\”
filepath = self.output_dir / filename
if filepath.exists():
return (ts_code, True, \”已存在\”)
# 获取数据
df = self.client.get_daily_data(ts_code, start_date, end_date)
if df is None or df.empty:
return (ts_code, False, \”无数据\”)
# 保存数据
df.to_csv(filepath, index=False)
return (ts_code, True, f\”成功 ({
len(df)}条)\”)
except Exception as e:
return (ts_code, False, str(e))
def fetch_batch(
self,
code_list: List[str],
start_date: str,
end_date: str,
show_progress: bool = True
) –> dict:
\”\”\”
批量获取数据
Args:
code_list: 股票代码列表
start_date: 开始日期
end_date: 结束日期
show_progress: 显示进度
Returns:
dict: 统计信息
\”\”\”
stats = {
\’total\’: len(code_list),
\’success\’: 0,
\’failed\’: 0,
\’skipped\’: 0,
\’errors\’: []
}
# 准备任务参数
tasks = [
(ts_code, start_date, end_date)
for ts_code in code_list
]
if show_progress:
print(f\”开始获取数据,共 {
len(code_list)} 支股票\”)
print(f\”日期范围:{
start_date} – {
end_date}\”)
print(f\”并发数:{
self.max_workers}\”)
print(\”-\” * 60)
start_time = time.time()
# 使用线程池
with ThreadPoolExecutor(max_workers=self.max_workers) as executor:
# 提交所有任务
future_to_code = {
executor.submit(self.fetch_single_stock, task): task[0]
for task in tasks
}
# 处理完成的任务
for i, future in enumerate(as_completed(future_to_code), 1):
ts_code, success, message = future.result()
if success:
stats[\’success\’] += 1
if message == \”已存在\”:
stats[\’skipped\’] += 1
symbol = \”⊘\”
else:
symbol = \”✓\”
else:
stats[\’failed\’] += 1
stats[\’errors\’].append((ts_code, message))
symbol = \”✗\”
if show_progress:
elapsed = time.time


