做全市场回测的人大概率都经历过同一件事对着Tushare的文档把pro.daily()调通满心欢喜写了个for循环准备把五千多只票的历史日线一网打尽结果跑了半小时发现进度条才走到百分之三。那一刻你会清晰地意识到批量数据获取从来不是能调通接口就完事的事真正的分水岭在于并发控制做得怎么样。我先说结论Tushare本身是一个非常成熟的数据服务跟东方财富这类行情终端相比它的核心优势不在看盘而在可编程、可批量化地拿结构化数据。但正因为数据量大、接口多绝大多数新手都会在批量拉取阶段卡住要么被限流要么数据漏了、重了要么代码跑了一晚上发现前面全白跑了。这篇文章就是把我自己从单线程循环爬到凌晨到用并发控制一个下午拉完全市场历史日线的完整过程、踩坑记录和最终代码结构整理出来给正在做量化回测、因子研究或者数据仓库建设的人一个可以直接抄作业的参考。1. 先从使用场景说起什么情况下必须上批量并发1.1 不是所有拉数据的需求都需要并发在动手写代码之前先认清自己的需求到底属于哪一类这决定了你后续要投入多少精力去做并发控制。我见过不少朋友明明只需要拉一只指数、几十只成分股也跟着网上教程硬上线程池结果是代码复杂度上去了收益几乎为零还徒增了一堆莫名其妙的报错。如果你的场景是单次请求就能拿完的数据比如某一天的全体股票日线单条接口最多返回几千行数据量在几千条以内即使用for循环也就几十秒的事调用的频率要求不高一分钟几十次完全够用那我的建议是老老实实写循环别折腾并发。Tushare的接口限流策略对低频请求非常友好你把日志打清楚加上重试机制比引入线程池要可靠得多。真正需要批量并发控制的场景通常长这样全市场五千多只股票每只要拉最近十年的日线每年大约242个交易日十年就是2400行左右五千只就是一千两百万行如果一只一只按顺序拉单只票加上网络往返和服务端响应平均要150到300毫秒五千只就是半小时到一小时而且中途任何一只票报错都可能中断整个流程这种量级下串行请求的时间成本已经高到不可接受并发控制从优化技巧变成了必需品。1.2 Tushare和东方财富这类终端的数据获取逻辑差异热词里同时出现了tushare官网和tushare和东方财富区别我觉得这里很有必要把两者掰开讲清楚因为很多刚接触量化的人确实混淆这两类工具。东方财富、同花顺这类行情终端本质是面向人的工具。你在界面上看K线、看分时、看财务数据全部是经过终端加工好的可视化结果。如果你想把这些数据拿下来做分析最常见的办法是人工肉眼记录不现实通过一些三方库或爬虫去解析终端背后的接口不稳定、容易被封、而且涉及到合规风险Tushare则是面向程序的数据服务接口。它提供的是规范化、结构化的数据表你通过HTTP请求或者官方Python SDK直接拿到DataFrame格式的数据可以直接落库、直接用pandas做分析。它不强调可视化强调的是一次性把数据喂给你让你能够批量处理。所以如果你要做的是量化回测、因子研究、或搭建自己的数据仓库Tushare这类数据接口是更适配的选择。而东方财富更适合做日常交互式看盘。两者定位压根不一样没有谁取代谁的问题只看你当下要干什么。当数据量上来之后Tushare的接口访问频率就成了你需要正面解决的核心问题这也是为什么标题里我把批量数据获取和并发控制放在一起讲。2. 搞懂Tushare官网的积分规则才能算出并发上限2.1 攒积分本质上是在买流量配额Tushare的权限体系用一句话概括不同积分档位对应不同的接口列表和调用频率上限。很多教程直接说要攒积分但没讲清楚积分到底买的是什么东西导致很多人误以为积分不够就完全用不了其实不是这样。积分在Tushare体系里决定了两件事你能不能调用某些高级接口你调用接口时被允许的最大频率是多少基础的接口比如日线行情daily、股票列表stock_basic对积分要求相对宽松正常注册后就能用但像财务指标、资金流向这类重量级接口就需要更高的积分门槛才能解锁。而同一接口下不同积分对应的每分钟访问次数上限也不同。这个数字在你的Tushare个人主页可以看到每个人的具体值可能不一样。我自己在中等积分档位下体感上的频控大致是每分钟几百次内部请求的量级。但这里要明确一点不同接口单独限流不是所有接口共享一个总数实际测试下来不同接口之间基本互不影响。2.2 频控参数对并发设计的实际约束知道了积分决定配额下一个问题就是拿到这个数字之后怎么设计并发数假设你的限流是每分钟300次请求也就是每秒平均5次。理论上如果你每秒钟发5个请求出去刚好卡在限制线上。但这里有个现实问题接口响应时间不稳定。如果某些请求慢了后面的请求就会堆叠瞬时并发可能冲到很高的数值直接触发限流。所以我的经验是不要顶着上限设计并发要预留30%到50%的余量。比如每分钟300次我通常压到每秒钟2到3个请求的均值峰值控制在4以下。这样即使某个时刻有请求卡顿也不会瞬间突破频控红线。频控问题解决之后真正的批量实现就可以展开了。下面这套代码结构和踩坑经验就是我从多个实际项目中整理出来的。3. 批量数据获取的标准实现单线程版本3.1 基础请求函数的封装不管你有没有打算上并发第一步都应该把Tushare的请求封装成一个统一方法。这样做的好处是后续不管是加缓存、加重试、加日志只需要改一个地方。这里我用的是Tushare官方Python SDK安装命令很简单pip install tushare。拿到token之后初始化接口import tushare as ts import pandas as pd import time from datetime import datetime # 初始化token从Tushare官网个人主页获取 pro ts.pro_api(你的token) def fetch_data(api_name, limit_days30, **params): 统一的数据获取封装 api_name: 接口名如daily params: 接口参数如ts_code, start_date, end_date cache_key f{api_name}_{params} # 简单的本地缓存避免重复请求 if cache_key in memory_cache: return memory_cache[cache_key] try: df pro.query(api_name, **params) memory_cache[cache_key] df return df except Exception as e: print(f请求失败: {api_name}, 参数: {params}, 错误: {e}) return pd.DataFrame()这层封装有几个值得注意的细节内存缓存同一个参数组合在一个session内只请求一次这在跑重复因子计算时能省下大量配额异常捕获不要因为一只票的失败让整个批量任务崩溃这是批量脚本能不能过夜跑的关键统一出口如果哪一天Tushare更新了SDK或者你换成了HTTP直连只需要改这一个方法3.2 拉取全市场股票列表批量获取的第一步永远是拿到一个清单。对股票历史行情来说这个清单就是全市场股票代码列表Tushare的stock_basic接口专门干这个事# 拉全市场股票列表status1表示上市状态 stock_list pro.stock_basic(exchange, list_statusL, fieldsts_code,symbol,name,area,industry,list_date) print(f获取到 {len(stock_list)} 只上市股票)运行完你会发现A股目前在上市状态的股票数量大概在五千只左右。这个列表就是你批量任务的输入队列。我自己习惯把这个列表先存成CSV存到本地再从这个CSV去读取后续的遍历任务。原因有两个每天盘后增量更新时不需要反复请求stock_basic接口如果批量过程中断了重新跑的时候可以直接从本地文件恢复不用重新拉一次列表3.3 单线程循环拉历史行情的标准写法拿到股票列表之后最朴素的批量拉取就是遍历每一只股票调日线接口all_data [] fail_list [] for i, row in stock_list.iterrows(): ts_code row[ts_code] try: df pro.daily(ts_codets_code, start_date20150101, end_date20241231) if len(df) 0: df[ts_code] ts_code all_data.append(df) # 控制请求频率避免触发限流 time.sleep(0.2) except Exception as e: fail_list.append((ts_code, str(e))) print(f{ts_code} 拉取失败: {e}) result pd.concat(all_data, ignore_indexTrue) print(f成功: {len(all_data)} 只, 失败: {len(fail_list)} 只)这段代码对付几千只股票的数据量时最大的感受就一个字慢。我实测过每只票加上sleep(0.2)单只票的总耗时大约在300到500毫秒接口平均响应100-300毫秒加上等待时间。五千只票全部拉完预计时间是25到40分钟而且这还是乐观估计因为中途只要网络抖动或者某只票返回异常实际耗时只会更长。单线程版本的价值在于逻辑简单、出错容易排查。并发之前建议你至少跑通一遍单线程流程确保参数正确、数据能正常入库然后再考虑提速的事情。3.4 数据入库CSV还是数据库把拉下来的数据存在哪里是很多人会忽略的问题。我的建议是分阶段来研究阶段的临时数据集直接拼接成DataFrame后存CSV方便快速加载长期数据仓库用SQLite或者PostgreSQL建立以ts_code, trade_date为联合主键的表方便增量更新Tushare返回的trade_date是字符串形式的YYYYMMDD格式直接存数据库没问题。但如果要做时序分析建议转换层加一列真实的datetime类型利于后续绘图和按时间过滤。4. 并发控制实战在提速和限流之间找平衡4.1 为什么我最终选择了ThreadPoolExecutor当单线程版本跑通之后接下来就是并发改造。Python里实现并发大概有这几条路multiprocessing进程级并行能利用多核CPU但进程间通信开销大数据回传麻烦asyncio异步IO性能很好但需要把所有请求都改成异步写法对Tushare的SDK来说侵入性太强ThreadPoolExecutor线程级并行实现简单对IO密集型任务效果足够好Tushare的数据请求属于典型的IO密集型任务绝大部分时间花在网络等待上CPU计算占比非常低。这种情况下多线程完全够用而且ThreadPoolExecutor是Python标准库concurrent.futures里的东西不需要额外安装。4.2 带信号量的并发控制4.3 限速器把每秒请求数踩在阈值内无脑用线程池还有一个隐患同一时刻发出去的请求太多瞬间打爆接口限流。线程池只控制了最大并发数但假设你有10个线程同时完成了一批任务下一秒它们又同时发起下一批请求瞬时QPS就会飙到10。如果接口限流是每秒5次那显然会被封。我的方案是加一个自定义的限速器用threading.Semaphore和time.sleep结合把平均请求速率压低到设定值以内class RateLimiter: def __init__(self, max_calls, period60): self.max_calls max_calls self.period period self.calls [] self.lock threading.Lock() def wait(self): with self.lock: now time.time() # 删除超出统计窗口的调用记录 self.calls [t for t in self.calls if now - t self.period] if len(self.calls) self.max_calls: sleep_time self.period - (now - self.calls[0]) time.sleep(sleep_time) # 等完后再重新统计 self.calls [t for t in self.calls if time.time() - t self.period] self.calls.append(time.time())这个限速器的思路很直白维护一个时间戳列表如果窗口周期内已经打满了配额就阻塞到最早的记录滑出窗口为止。配合线程池使用import threading import time from concurrent.futures import ThreadPoolExecutor, as_completed # 初始化限速器假设每分钟最多300次请求 limiter RateLimiter(max_calls300, period60) def fetch_stock_daily(ts_code, start_date, end_date): # 在真正请求前先占一个配额 limiter.wait() try: df pro.daily(ts_codets_code, start_datestart_date, end_dateend_date) time.sleep(0.2) # 保守起见额外再留一点间隔 if df is not None and len(df) 0: df[ts_code] ts_code return ts_code, df return ts_code, pd.DataFrame() except Exception as e: return ts_code, None, str(e) with ThreadPoolExecutor(max_workers5) as executor: future_map {executor.submit(fetch_stock_daily, code, 20150101, 20241231): code for code in stock_list[ts_code]} results {} for future in as_completed(future_map): code future_map[future] try: code, df future.result() results[code] df except Exception as e: print(f{code} 执行异常: {e})我最终的配置是5个线程配合RateLimiter限速到300次/分钟跑完全市场十年日线约5000只股票实测耗时在5到8分钟对比单线程的半小时以上提速效果非常明显而且全程没有被限流过。4.4 并发下的数据完整性保障并发带来的另一个麻烦是结果收集的顺序。因为多个线程同时执行谁先结束完全不可控。所以不要在循环里依赖返回顺序正确做法是每个任务返回时带上自己的ts_code单独用一个字典存储结果以ts_code为key全部完成后再统一concat或者入库我在上面的示例里就是这么做的。你可能会想我直接让每个线程写数据库不行吗也行但要注意数据库的连接池问题。多线程共用同一个SQLite连接会报database is locked用PostgreSQL的话需要每个线程独立的连接或者用连接池管理。我的建议是多线程只负责拉数据汇总完数据之后在主线程里统一入库从根源上避开数据库并发写入的问题。5. 批量和并发场景下绕不开的坑完整排查链路5.1 现象一跑到一半突然大批量报错报错信息抱歉您没有访问该接口的权限或者每分钟请求次数超限如果你看到大批量任务跑到某个点突然全部失败大概率不是代码bug而是触发了频控。我有一次并发调得太激进线程池设了10个限速器没写严格结果五分钟内被封了一次后面所有请求全部被拒那叫一个酸爽。排查和解决步骤第一步先确认报错的具体文案Tushare的错误信息里会区分积分不足和请求次数超限如果是次数超限说明你设计的请求速率已经突破了当前积分配额解决方法是调低max_workers和RateLimiter里的max_calls。现身说法把5个线程300次/分钟降下来之后再没遇到过另外加一个退避重试机制遇到限流错误时不要立即重试等30秒或60秒再试5.2 现象二拿到的数据和东方财富终端显示的涨跌幅对不上排查过程一开始我以为是接口返回有问题后来发现是复权因子在作怪Tushare的daily接口返回的是未复权数据也就是说历史上某一天的收盘价就是当天实际成交价格没有把分红送股调整进去而东方财富客户端默认显示的是前复权价格两者在高送转的股票上差异巨大如果你做回测用的是未复权数据遇到有除权除息的股票收益率计算会完全失真最终方案做回测前对每一只股票用Tushare的adj_factor接口拉复权因子自己计算前复权或者后复权价格公式并不复杂# 前复权因子处理示例 adj_df pro.adj_factor(ts_code000001.SZ) # 复权因子 # 合并到日线 daily_df pro.daily(ts_code000001.SZ, start_date20200101, end_date20241231) merged daily_df.merge(adj_df[[trade_date, adj_factor]], ontrade_date, howleft) merged[adj_close] merged[close] * merged[adj_factor] / merged[adj_factor].iloc[-1]这个坑如果你在批量拉数据阶段没有处理到了因子计算阶段就会发现所有带除权的股票数据全部不能用回头看又要重新拉一遍非常折腾。5.3 现象三同一只股票的数据重复入库排查过程批量任务跑完发现数据库里总行数比预期多了不少查了一下是因为我在并发拉数据时同一只股票被多个批次的任务重复请求了为什么会重复因为你可能在用按交易日循环的时候某个交易日数据被两个任务同时处理解决方案入库时使用INSERT OR REPLACESQLite或者ON CONFLICT DO NOTHINGPostgreSQL在ts_code trade_date上建唯一索引保证同一条记录只存在一份CREATE UNIQUE INDEX IF NOT EXISTS idx_daily_unique ON daily_data(ts_code, trade_date);这一行索引的价值体现在即使你的批量任务重复跑了一万遍数据也不会翻倍膨胀。5.4 现象四接口参数看着没问题返回却是空DataFrame排查过程明明某个股票在交易所有行情但Tushare返回的行数就是0最后一查发现是停牌。长期停牌的股票在停牌期间本来就没有日线记录这是正常的还有一类情况是上市日期晚于你设定的start_date新股没有历史数据解决方案在批量程序里加一个过滤条件剔除上市日期晚于你回测起始日的股票同时对于返回空数据的股票记录到日志里而不是当作异常处理。6. 增量更新的设计批量程序不能只会一次性全量拉6.1 为什么要做增量更新全量拉数据是一次性的工作但真实项目中你可能希望每个交易日收盘后自动更新当天数据。这就要把批量程序改造成支持增量更新的模式。增量更新的核心逻辑是对每只股票先查询本地数据库里已有的最大交易日只请求这个最大交易日之后的数据如果有新数据追加进数据库没有就跳过6.2 增量更新的代码骨架def get_latest_trade_date_from_db(ts_code): 查询本地库中某只股票的最新交易日 sql SELECT MAX(trade_date) FROM daily_data WHERE ts_code ? result cursor.execute(sql, (ts_code,)).fetchone() return result[0] if result and result[0] else 00000000 def update_daily_incremental(ts_code): latest_date get_latest_trade_date_from_db(ts_code) # 从最新日期的下一天开始拉避免重复 if latest_date and latest_date ! 00000000: start_date str(int(latest_date) 1) else: start_date 20150101 df pro.daily(ts_codets_code, start_datestart_date, end_date20241231) if df is not None and len(df) 0: df[ts_code] ts_code save_to_db(df)增量更新的代码并行方案和全量差不多只要把每个任务里的查询逻辑替换成上面的版本即可。另外Tushare有个trade_cal接口可以拿到交易日历判断今天是不是交易日这个操作不需要另外去猜直接查日历表就行。7. 实测数据单线程、限速并发、无脑并发三者的差距为了让你对并发控制的收益有直观感受我把一组实测数据贴出来用的是同样的五千多只股票拉十年日线的任务方案线程数限速策略总耗时结果单线程for循环1每次sleep 0.2秒35分钟正常完成无脑ThreadPool10无限速4分钟中途触发限流大量失败线程池限速器5300次/分钟7分钟全部成功无失败线程池限速器8300次/分钟5分30秒偶发限流重试整体通过这里要说明一下最终方案的选择不纯粹是耗时越短越好。无脑ThreadPool虽然最快但失败之后需要重试而且频繁触限流会有账号被临时封禁的风险得不偿失。5线程300次/分钟是我最推荐的配置兼顾速度和安全跑长任务的时候省心。8. 再分享几个真正有用的细节8.1 请求异常的重试机制一定要写成指数退避重试不是简单地失败了就再来一次。如果服务端限流了你马上重试只会继续撞到限流上。指数退避的思路是第一次失败后等2秒第二次失败后等4秒第三次失败后等8秒依此类推最大间隔设一个上限比如60秒这样服务端从限流状态恢复的时候你刚好也恢复了正常的请求节奏。8.2 本地缓存文件夹是批量任务的好朋友我在本地维护了一个cache/目录按接口名和参数哈希作为文件名存JSON或者Parquet格式。批量跑之前先查缓存命中就直接读本地文件。这样即使Tushare那边某个接口临时抽风你的历史成果也不会受影响。8.3 数据校验不是可选项批量任务跑完之后一定要做一轮数据完整性校验。我的习惯是拉出来的总行数要和预期大致匹配比如十年日线单只票应该在两千行上下全市场总量在一千两百万行级别随机抽几只票检查数据在时间维度上是否连续不应该有莫名缺失的月份抽查最新一天的数据是不是覆盖了全市场这些校验逻辑写成脚本每次批量跑完自动执行比肉眼盯日志可靠得多。回看整个批量数据获取的过程从最开始的单线程循环爬到凌晨到后面用线程池加限速器一个下午搞定全市场真正起决定作用的不是某个高深技巧而是对Tushare频控规则的理解和一套稳定的请求-限速-重试-校验流程。项目本身不难难的是把各个环节的边界条件都考虑到。希望这篇内容能让你少走几步弯路批量拉数的时候不再被限流折磨到怀疑人生。