本文是「Python量化实战」系列第 16 篇,从这篇起进入工程化模块:前 15 篇解决"数据怎么取、策略怎么写",接下来 5 篇解决"怎么把它做成一个稳定运转的系统"。

写量化脚本的人几乎都经历过这个阶段:单只股票的数据拉取调通了,兴冲冲地套个 for 循环拉全市场——然后发现 5000 多只股票要跑将近一个小时,中途网络抖一下还前功尽弃,只能从头再来。批量拉取从来不是"加个循环"那么简单,它是一个标准的工程问题:并发提速、失败重试、进度可见、结果合并,四件事缺一不可。本文用不到 100 行代码把这四件事一次做对,实测并发方案比串行提速 6.9 倍,且 40 只样本全部成功、0 失败。

本文你将得到什么

  1. 一套线程池并发拉取模板(ThreadPoolExecutor + as_completed,可直接套用到任何接口)
  2. 指数退避的失败重试装饰逻辑(网络抖动不再毁掉整批任务)
  3. 实时进度输出 + 失败清单(跑到哪、挂了谁,一目了然)
  4. 串行 vs 并发的真实耗时对比数据(0.46 秒/只 → 0.07 秒/只)
  5. 线程数怎么定、限频怎么躲的工程经验

一、在线体验

想先在线试试接口效果?打开 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
  • 关注公众号获取本系列更新

本文为技术演示,不构成投资建议。

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐