【魔码量化工程实战进阶 #01】拉数管道设计:从一次性脚本到可重启的数据流水线
摘要:【魔码量化工程实战进阶 #01】拉数管道设计:从"一次性脚本"到"可重启的数据流水线 入门系列教你怎么"拉一只",本篇教你"稳稳拉完 5000 只还不怕中途崩"。这是整个
【魔码量化工程实战进阶 #01】拉数管道设计:从"一次性脚本"到"可重启的数据流水线
入门系列教你怎么"拉一只",本篇教你"稳稳拉完 5000 只还不怕中途崩"。这是整个工程体系的地基——后面所有模块都建立在一个"能断点续拉、失败不垮"的管道之上。
本文你将得到什么
- 为什么
for code in codes:这种一次性脚本在生产环境必翻车 - 数据流水线的四个工程属性:幂等 / 断点续拉 / 状态持久化 / 失败隔离
- 一个可直接套用的
Puller骨架(checkpoint + 续拉 + 单只容错) - 五个真实踩过的坑(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 只全市场来说,这意味着一次中断最多损失"中断后还没拉的那几只",而不是全部。
五、原理深挖:五个真实踩过的坑
- checkpoint 写时机错:有人等"全部拉完"才写
state.json。一旦中途崩,进度全丢,等于没做。正确做法:拉完一只、落库成功、立即写 checkpoint。 - 状态文件并发冲突:多进程/多线程共享同一个
state.json,会出现"读旧→写旧"覆盖。解决:各进程分片(进程 A 管 000–999,进程 B 管 1000–1999),或用文件锁。 - 幂等键不全:只用
code当唯一键,忘了加"区间/日期"。结果断点续拉时,昨天拉的今天被当成"已完成"跳过,漏掉增量。唯一键必须是(code, 区间)。 - 额度浪费:没有续拉逻辑时,每次重跑都全量拉 5000 只。有了 checkpoint,重跑只补缺失,免费证额度/限频配额省下 90%+。
- 死信无限重试:单只连续失败还死循环重试,把限频配额打光、把任务卡死。正确做法:连续失败 N 次进死信队列,标记为"人工核查",绝不无限重试。
六、小结
数据管道的核心不是"会不会发请求",而是"崩了能不能恢复、重跑会不会浪费"。一个 state.json + 单只容错,就把脆弱的脚本变成了可重启的流水线。这是后面所有模块(并发、限频、增量、调度)共同依赖的地基。
全市场 5000+ 只一次拉完,免费证的额度与限频往往不够用——这正是魔码量化 Pro 包(更高配额、完整历史)存在的意义。详见文末。
免责声明:本文所有示例数据仅用于接口演示,不构成任何投资建议;市场有风险,投资需谨慎。
系列持续更新中。 想要亲手跑通上面的代码?前往 魔码证书申请页 免费领取你的专属证书,复制即用、按次计费、稳定可用。
