【Python量化实战 #16】全市场5000只股一只只拉要等到哭?Python并发批量拉取工程方案(重试+进度+合并)
本文是「Python量化实战」系列第 16 篇,从这篇起进入工程化模块:前 15 篇解决"数据怎么取、策略怎么写",接下来 5 篇解决"怎么把它做成一个稳定运转的系统"。
写量化脚本的人几乎都经历过这个阶段:单只股票的数据拉取调通了,兴冲冲地套个 for 循环拉全市场——然后发现 5000 多只股票要跑将近一个小时,中途网络抖一下还前功尽弃,只能从头再来。批量拉取从来不是"加个循环"那么简单,它是一个标准的工程问题:并发提速、失败重试、进度可见、结果合并,四件事缺一不可。本文用不到 100 行代码把这四件事一次做对,实测并发方案比串行提速 6.9 倍,且 40 只样本全部成功、0 失败。
本文你将得到什么
- 一套线程池并发拉取模板(
ThreadPoolExecutor+as_completed,可直接套用到任何接口) - 带指数退避的失败重试装饰逻辑(网络抖动不再毁掉整批任务)
- 实时进度输出 + 失败清单(跑到哪、挂了谁,一目了然)
- 串行 vs 并发的真实耗时对比数据(0.46 秒/只 → 0.07 秒/只)
- 线程数怎么定、限频怎么躲的工程经验
一、在线体验
想先在线试试接口效果?打开 API Playground 即可直接调用测试:
https://mairuiapi.com/playground
本文用到两个接口:全市场股票列表(stock_list)和历史 K 线(stock_history),可以先在 Playground 里看看返回结构。
二、环境准备
本文代码使用 mairui SDK 获取股票数据,安装方法如下:
pip install mairui
SDK 的完整接口文档与使用说明请查阅 GitHub 仓库:https://github.com/MaiRuiApi/mairui
接口的详细参数说明请查阅官网 API 文档:https://mairuiapi.com/hsdata
运行环境:
- Python 3.9+(
concurrent.futures为标准库,无需额外安装) - mairui SDK 1.0.0
- pandas 2.2+
import os
import mairui
# 证书从环境变量读取(官网注册后获取),不要写死在代码里
api = mairui.Client("LICENCE-66D8-9F96-0C7F0FBCD073") # 证书从环境变量读取
本文数据截至 2026-07-28,拉取区间 2026-06-01 至 2026-07-28,读者复现时数据可能略有差异。
三、先拿到"任务清单":全市场股票列表
批量任务的第一步是确定任务边界。stock_list 一次返回全市场 A 股列表:
stock_list = api.stock_list()
print(f"全市场 A 股数量: {len(stock_list)}")
# 每个元素形如 {"dm": "000001.SZ", "mc": "平安银行", "jys": "SZ"}
真实运行输出:
全市场 A 股数量: 5205
5205 只股票就是我们的任务全集。本文演示取前 40 只(把样本换成全量列表,代码完全不用改,只是跑得久一点)。
四、单只拉取函数:先把"重试"做进去
并发的前提是单任务函数足够健壮。网络请求天然会失败——超时、抖动、瞬时限频——正确姿势不是祈祷不失败,而是失败了自动重试,并且一次比一次等得久(指数退避):
import time
import pandas as pd
MAX_RETRIES = 3 # 单只失败最大重试次数
RETRY_BACKOFF = 1.5 # 退避基数:第 k 次重试等待 1.5^k 秒
ST, ET = "20260601", "20260728" # 拉取区间(st/et 格式为 YYYYMMDD,不带横杠)
def fetch_one(api, code, name):
"""拉单只股票日K,带重试与指数退避。失败到底则抛出最后一次异常。"""
last_exc = None
for attempt in range(MAX_RETRIES):
try:
kline = api.stock_history(code.split(".")[0], "d", "n", st=ST, et=ET)
kline_df = pd.DataFrame(kline)
kline_df.insert(0, "code", code) # 打上股票标识,合并后可区分
kline_df.insert(1, "name", name)
return kline_df
except Exception as exc:
last_exc = exc
time.sleep(RETRY_BACKOFF ** (attempt + 1)) # 1.5s → 2.25s → 3.4s
raise last_exc
关键解释:
- 每个 DataFrame 拉下来立刻
insert股票代码列——合并之后还能知道每行属于谁,这是新手最常漏的一步; - 指数退避的意义:如果失败是限频导致的,立刻重试只会继续被拒,等待时间递增才能自愈。
五、并发拉取:线程池 + as_completed
数据拉取是典型的 IO 密集型任务(时间都花在等网络响应上),线程池就是标准答案:
from concurrent.futures import ThreadPoolExecutor, as_completed
MAX_WORKERS = 5 # 线程数控制在 5 以内,避免触发接口频率限制
def run_concurrent(api, stocks):
"""并发拉取:进度实时输出 + 失败清单登记。"""
frames, failed = [], []
with ThreadPoolExecutor(max_workers=MAX_WORKERS) as pool:
futures = {pool.submit(fetch_one, api, s["dm"], s["mc"]): s for s in stocks}
done = 0
for fut in as_completed(futures): # 谁先完成先处理谁
stock = futures[fut]
done += 1
try:
frames.append(fut.result())
except Exception as exc: # 重试后仍失败:登记,不中断整批
failed.append((stock["dm"], type(exc).__name__))
if done % 10 == 0 or done == len(stocks):
print(f" 进度 {done}/{len(stocks)},失败 {len(failed)}")
merged = pd.concat(frames, ignore_index=True) if frames else pd.DataFrame()
return merged, failed
关键解释:
as_completed按完成顺序返回结果,天然适合做进度条;- 单只失败绝不抛出中断整批,而是记入
failed清单,批量任务跑完后单独补拉失败部分——这是批量工程的基本素养。
六、真实耗时对比:串行 0.46 秒/只 vs 并发 0.07 秒/只
同一台机器、同一网络环境的真实测试(2026-07-29 验证):
Step2 串行基准(前 10 只)
串行拉取 10 只耗时: 4.6 秒(约 0.46 秒/只)
Step3 并发拉取(40 只,5 线程,重试上限 3)
进度 10/40,失败 0
进度 20/40,失败 0
进度 30/40,失败 0
进度 40/40,失败 0
并发拉取 40 只耗时: 2.7 秒(约 0.07 秒/只)
提速比(按单只均摊): 6.9x
失败清单: 无
按此速度外推:全市场 5205 只,串行约 40 分钟,5 线程并发约 6 分钟。而且这 6 分钟是"带重试保险"的 6 分钟——中途抖动自动兜住,不再需要从头重跑。
合并结果同样一步到位:
合并后 DataFrame: 1640 行 x 11 列,覆盖 40 只股票
code name a c h l o pc sf t v
000001.SZ 平安银行 1042306455.0 10.99 10.99 10.81 10.90 10.93 0 2026-06-01 954596
000001.SZ 平安银行 978159336.0 11.08 11.10 10.94 10.98 10.99 0 2026-06-02 885428
七、完整可运行示例
# -*- coding: utf-8 -*-
import os
import time
from concurrent.futures import ThreadPoolExecutor, as_completed
import pandas as pd
import mairui
SAMPLE_SIZE = 40
MAX_WORKERS = 5
MAX_RETRIES = 3
RETRY_BACKOFF = 1.5
ST, ET = "20260601", "20260728"
# fetch_one / run_concurrent 定义见上文第四、五节
def main():
api = mairui.Client("LICENCE-66D8-9F96-0C7F0FBCD073") # 证书从环境变量读取
stock_list = api.stock_list() # ① 任务清单
sample = stock_list[:SAMPLE_SIZE] # 全市场就把切片去掉
merged_df, failed = run_concurrent(api, sample) # ② 并发拉取
print(f"成功 {merged_df['code'].nunique()} 只,失败清单: {failed or '无'}")
merged_df.to_csv("batch_kline.csv", index=False, encoding="utf-8-sig") # ③ 落盘
if __name__ == "__main__":
main()
八、避坑与进阶
- 坑 1:线程开得越多越快? 不是。接口侧有频率限制,线程数超过阈值后失败率飙升,重试反而拖慢整体。实测 3~5 线程是稳定与速度的平衡点。
- 坑 2:用多进程做 IO 任务。
multiprocessing适合 CPU 密集型计算;拉数据是 IO 等待,线程池更轻、共享内存更方便。 - 坑 3:合并时索引错乱。
pd.concat记得ignore_index=True,否则各 DataFrame 的行索引会重复。 - 坑 4:失败任务无记录。批量任务必须产出"失败清单",跑完针对性补拉,而不是整批重来。
- 进阶方向:拉下来的数据往哪存?CSV 会越来越慢——下一篇讲 CSV/SQLite/MySQL 三种存储方案的对比与选型,让 5000 只股票的历史数据"存得下、查得快、能增量"。
九、总结与延伸
批量拉取的工程要点浓缩成一句话:并发提速、重试兜底、进度可见、失败可补。这套模板不只适用于 K 线——财务数据、实时快照、公告列表,任何"按代码逐只拉取"的场景都能直接套用。配合稳定的数据接口,全市场级的数据任务从"跑一次要祈祷"变成"每天定时无人值守"。
延伸阅读:
- 在线体验更多接口:https://mairuiapi.com/playground
- 查看完整 API 文档:https://mairuiapi.com/hsdata
- SDK 文档与源码:https://github.com/MaiRuiApi/mairui
- 关注公众号获取本系列更新
本文为技术演示,不构成投资建议。
更多推荐


所有评论(0)