Tushare 实战:批量数据获取与并发控制
批量数据获取与并发控制前一篇完成了环境搭建和单个股票数据获取。这篇解决实际问题:如何高效获取全市场 5000+ 支股票的数据?问题分析单线程获取的瓶颈# 单线程获取forts_codeinstock_list:df=client.get_daily_data(ts_code,'20230101','20231231')time.sleep(0.1)# 频率控制# 问题:# - 5000 支股票 × 0.1 秒 = 500 秒 ≈ 8.3 分钟# - 如果获取 5 年数据,需要 40+ 分钟# - 网络异常会导致整个流程中断需要解决的问题速度慢:单线程效率低易中断:网络异常导致前功尽弃频率限制:触发 API 限流进度丢失:无法断点续传资源浪费:CPU 和网络空闲解决方案:并发控制方案对比方案优点缺点适用场景多线程简单,IO 友好GIL 限制IO 密集型 ✓多进程突破 GIL内存占用大CPU 密集型异步 IO性能最好学习成本高高并发场景进程池简单易用灵活性一般中等规模 ✓选择:进程池(ProcessPoolExecutor)+ 线程池(ThreadPoolExecutor)组合实现方案方案 1:线程池(推荐入门)创建src/batch_fetch_thread.py:"""批量获取日线数据 - 线程池版本"""fromdata_clientimportTushareClientfromconcurrent.futuresimportThreadPoolExecutor,as_completedfrompathlibimportPathimporttimefromtypingimportList,TupleclassBatchFetcher:"""批量数据获取器"""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)deffetch_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=argstry:# 检查是否已存在filename=f"{ts_code.replace('.','_')}.csv"filepath=self.output_dir/filenameiffilepath.exists():return(ts_code,True,"已存在")# 获取数据df=self.client.get_daily_data(ts_code,start_date,end_date)ifdfisNoneordf.empty:return(ts_code,False,"无数据")# 保存数据df.to_csv(filepath,index=False)return(ts_code,True,f"成功 ({len(df)}条)")exceptExceptionase:return(ts_code,False,str(e))deffetch_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)forts_codeincode_list]ifshow_progress:print(f"开始获取数据,共{len(code_list)}支股票")print(f"日期范围:{start_date}-{end_date}")print(f"并发数:{self.max_workers}")print("-"*60)start_time=time.time()# 使用线程池withThreadPoolExecutor(max_workers=self.max_workers)asexecutor:# 提交所有任务future_to_code={executor.submit(self.fetch_single_stock,task):task[0]fortaskintasks}# 处理完成的任务fori,futureinenumerate(as_completed(future_to_code),1):ts_code,success,message=future.result()ifsuccess:stats['success']+=1ifmessage=="已存在":stats['skipped']+=1symbol="⊘"else:symbol="✓"else:stats['failed']+=1stats['errors'].append((ts_code,message))symbol="✗"ifshow_progress:elapsed=time.time