← 返回博客列表

【魔码量化工程实战进阶 #01】拉数管道设计:从一次性脚本到可重启的数据流水线

2026年08月31日 17:59 · 魔码数服 · 魔码量化工程实战进阶

摘要:【魔码量化工程实战进阶 #01】拉数管道设计:从"一次性脚本"到"可重启的数据流水线 入门系列教你怎么"拉一只",本篇教你"稳稳拉完 5000 只还不怕中途崩"。这是整个

【魔码量化工程实战进阶 #01】拉数管道设计:从"一次性脚本"到"可重启的数据流水线

入门系列教你怎么"拉一只",本篇教你"稳稳拉完 5000 只还不怕中途崩"。这是整个工程体系的地基——后面所有模块都建立在一个"能断点续拉、失败不垮"的管道之上。

本文你将得到什么

  1. 为什么 for code in codes: 这种一次性脚本在生产环境必翻车
  2. 数据流水线的四个工程属性:幂等 / 断点续拉 / 状态持久化 / 失败隔离
  3. 一个可直接套用的 Puller 骨架(checkpoint + 续拉 + 单只容错)
  4. 五个真实踩过的坑(checkpoint 写时机、状态文件并发、额度浪费…)

一、痛点:一次性脚本的两种死法

新手拉数常写成这样:

for code in all_codes:          # all_codes 可能是 5000+ 只
    data = get(code)            # 第 2333 只时网络抖了一下
    save(data)

它有两种死法:

  • 死法 A(中途崩):第 2333 只请求超时,异常向上抛,前面 2332 只白拉。重跑又从第 0 只开始——免费证额度、限频配额被重复消耗。
  • 死法 B(盘中禁跑):A 股 9:30–11:30、13:00–15:00 严禁任何重跑/补数(运维红线)。一旦盘中脚本崩了,你只能干等收盘,当天数据缺口补不上。

本质问题:脚本没有"记忆"。它不记得自己拉到哪了,也不隔离单只失败。


二、工程方案:把"脚本"升级成"流水线"

给拉数加上四个属性,它就从"一次性脚本"变成了"可重启的数据流水线":

属性 含义 怎么实现
幂等 同一只、同一区间拉两次,结果一致、不重复写 落库用 (code, 区间) 唯一键,INSERT OR IGNORE
断点续拉 进程被杀后重启,从断点继续 用状态文件 state.json 记录已完成的 code
状态持久化 进度落盘,不依赖内存 checkpoint 每拉完一只就写盘
失败隔离 一只挂了不阻断整体 单只 try/except,失败记死信列表,最后统一重试

三、可跑代码(零 SDK,纯 HTTP;演示证书标注"演示数据")

下面这个 Puller 是一个生产可用的骨架。演示证书返回的是演示数据,真实使用把 TOKEN 换成你的正式证书即可。

import json, os, time, requests

TOKEN = "TEST-API-TOKEN-MOMA-836089C22111"   # 换成你的魔码正式证书
BASE  = "https://api.momaapi.com"
STATE = "pull_state.json"                    # checkpoint 文件

class Puller:
    def __init__(self, state_file=STATE):
        self.state_file = state_file
        self.done = self._load()             # 已完成的 code 集合

    def _load(self):
        if os.path.exists(self.state_file):
            return set(json.load(open(self.state_file, encoding="utf-8")).get("done", []))
        return set()

    def _save(self):
        # 每拉完一只就写盘:进度不依赖内存,进程被杀也不丢
        json.dump({"done": sorted(self.done)}, open(self.state_file, "w", encoding="utf-8"))

    def fetch(self, code):
        # 真实场景:拉历史/实时;这里用实时快照示意
        r = requests.get(f"{BASE}/hsstock/real/time/{code}/{TOKEN}", timeout=20)
        r.raise_for_status()
        return r.json()

    def run(self, codes, max_retry=3):
        dead = []                             # 死信列表:连续失败的单只
        for code in codes:
            if code in self.done:            # 断点续拉:跳过已完成的
                continue
            ok = False
            for _ in range(max_retry):       # 单只失败隔离:重试不阻断整体
                try:
                    self.fetch(code)
                    ok = True
                    break
                except Exception as e:
                    time.sleep(0.3)
            if ok:
                self.done.add(code)
                self._save()                 # checkpoint 写时机:拉完即写
            else:
                dead.append(code)
        return dead

# 用法
codes = ["600519", "000001.SZ", "300750", "601318", "000858",
         "002594", "600036", "601012", "600276", "000333"]
dead = Puller().run(codes)
print("未拉取成功(死信):", dead)

关键点:① if code in self.done: continue 实现断点续拉;② _save() 在每只成功后立即落盘,进程中途被 kill -9 也不丢进度;③ 单只失败被 try/except 吞掉,记进 dead,不阻断后面的 code。


四、本机实测(本地模拟列表演示续拉逻辑;真实场景替换 mock 为 API)

为了不消耗接口额度,下面用本地模拟列表演示"中途打断 → 重启续拉只补缺失"的行为(逻辑与真实 API 完全一致):

import json, os
STATE = "demo_state.json"
MOCK = [f"C{i:04d}" for i in range(10)]     # 模拟 10 只待拉

class DemoPuller:
    def __init__(s):
        s.done = set(json.load(open(STATE, encoding="utf-8")).get("done", [])) if os.path.exists(STATE) else set()
    def _save(s):
        json.dump({"done": sorted(s.done)}, open(STATE, "w", encoding="utf-8"))
    def run(s, codes):
        for c in codes:
            if c in s.done: continue
            s.done.add(c); s._save()         # 模拟"拉取+落库+写checkpoint"
        return s.done

# 第一次:只跑前 5 只,然后"进程被杀"
open(STATE, "w").write('{"done":[]}')
p = DemoPuller()
p.run(MOCK[:5])
print("第一次跑完,state =", sorted(p.done))   # ['C0000'..'C0004']

# 第二次:重启,跑全部 10 只 —— 只补 C0005..C0009
p2 = DemoPuller()
p2.run(MOCK)
print("重启续拉后,state =", sorted(p2.done))  # 全部 10 只,前 5 只没重拉
print("补拉数量 =", len(p2.done) - 5, "只(只补缺失,不浪费配额)")

真实输出:

第一次跑完,state = ['C0000', 'C0001', 'C0002', 'C0003', 'C0004']
重启续拉后,state = ['C0000', 'C0001', 'C0002', 'C0003', 'C0004', 'C0005', 'C0006', 'C0007', 'C0008', 'C0009']
补拉数量 = 5 只(只补缺失,不浪费配额)

这就是断点续拉的价值:重启不重拉已完成的,只补缺口。对 5000 只全市场来说,这意味着一次中断最多损失"中断后还没拉的那几只",而不是全部。


五、原理深挖:五个真实踩过的坑

  1. checkpoint 写时机错:有人等"全部拉完"才写 state.json。一旦中途崩,进度全丢,等于没做。正确做法:拉完一只、落库成功、立即写 checkpoint
  2. 状态文件并发冲突:多进程/多线程共享同一个 state.json,会出现"读旧→写旧"覆盖。解决:各进程分片(进程 A 管 000–999,进程 B 管 1000–1999),或用文件锁。
  3. 幂等键不全:只用 code 当唯一键,忘了加"区间/日期"。结果断点续拉时,昨天拉的今天被当成"已完成"跳过,漏掉增量。唯一键必须是 (code, 区间)
  4. 额度浪费:没有续拉逻辑时,每次重跑都全量拉 5000 只。有了 checkpoint,重跑只补缺失,免费证额度/限频配额省下 90%+。
  5. 死信无限重试:单只连续失败还死循环重试,把限频配额打光、把任务卡死。正确做法:连续失败 N 次进死信队列,标记为"人工核查",绝不无限重试。

六、小结

数据管道的核心不是"会不会发请求",而是"崩了能不能恢复、重跑会不会浪费"。一个 state.json + 单只容错,就把脆弱的脚本变成了可重启的流水线。这是后面所有模块(并发、限频、增量、调度)共同依赖的地基。

全市场 5000+ 只一次拉完,免费证的额度与限频往往不够用——这正是魔码量化 Pro 包(更高配额、完整历史)存在的意义。详见文末。


免责声明:本文所有示例数据仅用于接口演示,不构成任何投资建议;市场有风险,投资需谨慎。

系列持续更新中。 想要亲手跑通上面的代码?前往 魔码证书申请页 免费领取你的专属证书,复制即用、按次计费、稳定可用。

想亲自试一下?免费获取证书
客服微信
客服微信二维码