Tushare批量获取全市场数据:并发控制实战与踩坑指南
2026/9/24 19:23:30 网站建设 项目流程

做全市场回测的人,大概率都经历过同一件事:对着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_days=30, **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接口专门干这个事:

# 拉全市场股票列表,status=1表示上市状态 stock_list = pro.stock_basic(exchange='', list_status='L', fields='ts_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_code=ts_code, start_date='20150101', end_date='20241231') 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_index=True) 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.Semaphoretime.sleep结合,把平均请求速率压低到设定值以内:

class RateLimiter: def __init__(self, max_calls, period=60): 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_calls=300, period=60) def fetch_stock_daily(ts_code, start_date, end_date): # 在真正请求前先占一个配额 limiter.wait() try: df = pro.daily(ts_code=ts_code, start_date=start_date, end_date=end_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_workers=5) 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_workersRateLimiter里的max_calls。现身说法,把5个线程+300次/分钟降下来之后,再没遇到过
  • 另外加一个退避重试机制:遇到限流错误时,不要立即重试,等30秒或60秒再试

5.2 现象二:拿到的数据和东方财富终端显示的涨跌幅对不上

排查过程

  • 一开始我以为是接口返回有问题,后来发现是复权因子在作怪
  • Tushare的daily接口返回的是未复权数据,也就是说历史上某一天的收盘价就是当天实际成交价格,没有把分红送股调整进去
  • 而东方财富客户端默认显示的是前复权价格,两者在高送转的股票上差异巨大
  • 如果你做回测用的是未复权数据,遇到有除权除息的股票,收益率计算会完全失真

最终方案:做回测前,对每一只股票用Tushare的adj_factor接口拉复权因子,自己计算前复权或者后复权价格,公式并不复杂:

# 前复权因子处理示例 adj_df = pro.adj_factor(ts_code='000001.SZ') # 复权因子 # 合并到日线 daily_df = pro.daily(ts_code='000001.SZ', start_date='20200101', end_date='20241231') merged = daily_df.merge(adj_df[['trade_date', 'adj_factor']], on='trade_date', how='left') merged['adj_close'] = merged['close'] * merged['adj_factor'] / merged['adj_factor'].iloc[-1]

这个坑如果你在批量拉数据阶段没有处理,到了因子计算阶段就会发现所有带除权的股票数据全部不能用,回头看又要重新拉一遍,非常折腾。

5.3 现象三:同一只股票的数据重复入库

排查过程

  • 批量任务跑完,发现数据库里总行数比预期多了不少
  • 查了一下,是因为我在并发拉数据时,同一只股票被多个批次的任务重复请求了
  • 为什么会重复?因为你可能在用"按交易日循环"的时候,某个交易日数据被两个任务同时处理

解决方案:入库时使用INSERT OR REPLACE(SQLite)或者ON CONFLICT DO NOTHING(PostgreSQL),在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_code=ts_code, start_date=start_date, end_date='20241231') 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频控规则的理解和一套稳定的"请求-限速-重试-校验"流程。项目本身不难,难的是把各个环节的边界条件都考虑到。希望这篇内容能让你少走几步弯路,批量拉数的时候不再被限流折磨到怀疑人生。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询