本节摘要:行情是量化系统的心跳。本节先分通道:REST 适合拉历史与低频快照,WebSocket 适合订阅实时推送,一张对比表给出选型依据;再分数据:K 线(研究主力)、深度(执行视角)、成交(微观结构)三类数据的形态、用途与更新节奏;然后讲接入的四个工程要点——限速合规、断线重连、心跳与时间校准、闭合为准的落库纪律(data_sources 为官方模块,行为细节以官方文档为准);最后用一个最小消费骨架示意数据从接口到规范化字典的旅程,并以一张新数据源接入走查清单与高频排查对照收尾。本节不追求覆盖所有交易所,只建立一套可以平移到任何数据源的接入方法论。
| 维度 | REST 拉取 | WebSocket 订阅 |
|---|---|---|
| 交互模式 | 一问一答 | 建连后持续推送 |
| 擅长 | 历史 K 线回填、低频快照、对账核对 | 实时成交、深度增量、盘口变化 |
| 成本特征 | 受限速配额约束,重拉昂贵 | 长连接维护成本,断线需补数 |
| 典型用途 | 初始化历史、3.3 节缺口回填 | trading-worker 驱动策略的实时输入 |
经验法则:问过去用 REST,看现在用 WebSocket;两者不是二选一,而是初始化与运行的接力——首启用 REST 建历史,运行期用 WebSocket 保新鲜,断线后用 REST 补洞。
| 数据 | 形态 | 主要用途 | 更新节奏 |
|---|---|---|---|
| K 线 | 开高低收 + 成交量,按周期聚合 | 指标与策略的主战场(第 4、5 章)、回测(第 6 章) | 每周期闭合一根 |
| 深度 | 按价位分档的挂单量 | 执行质量分析、大额订单的冲击预估(示意用途) | 流式快照/增量 |
| 成交 | 逐笔或聚合成交记录 | 微观结构研究、高频视角的验证 | 实时流 |
三类数据的价值密度与存储成本同向递增:K 线最省最常用,逐笔成交最贵最专用。个人自托管的合理起点是 K 线为主、深度按需、成交谨慎——存储预算(第 1.1 节的 20 GB 量级为示意值)会替你做这个决定。
**要点一:限速合规。**交易所对 REST 与订阅都有配额(以各交易所官方文档为准)。合规姿势:客户端主动节流(令牌桶一类),把重试做成带退避的指数间隔,429 类响应视为「停一停」信号而不是「再试一次」信号。
限速最常用的实现是令牌桶(教学示意):
# rate_limit_demo.py —— 令牌桶限速的最小示意(纯标准库) import time class TokenBucket: def __init__(self, capacity, refill_per_sec): self.capacity = capacity # 桶容量(示意值) self.tokens = capacity self.refill = refill_per_sec # 每秒补充速率(示意值) self.last = time.monotonic() def allow(self, n=1): now = time.monotonic() self.tokens = min(self.capacity, self.tokens + (now - self.last) * self.refill) self.last = now if self.tokens >= n: self.tokens -= n return True return False # 拿不到令牌就等,而不是硬闯 if __name__ == "__main__": bucket = TokenBucket(capacity=5, refill_per_sec=2) results = [bucket.allow() for _ in range(8)] print(results) # 示意:前若干个立即通过,其余按补充速率排队
桶的两个参数对应两种意图:容量容纳突发(比如重连后的一次性补数),速率守住长期配额。把它挂在所有外呼之前,比在收到限速响应后手忙脚乱地退避体面得多——后者是补救,前者是纪律。
**要点二:断线重连。**长连接必然断。重连流程:退避重连 ──▶ 重新订阅 ──▶ 用 REST 补齐断线窗口 ──▶ 恢复消费。断线窗口补不齐就形成缺口,进入 3.3 节质检流程。
**要点三:心跳与时间校准。**维持连接健康靠心跳;对齐「现在」靠时间源——本地时钟漂移会让「这根 K 线属于哪一分钟」都变成问题。约定:对内一律 UTC 存储,展示层再转本地时区(3.3 节展开)。
**要点四:闭合为准。**一根 K 线只有走完才是事实;未闭合的当前根是「进行时」。落库纪律:闭合根入正式表供研究与回测,进行时根只供展示。这条纪律是第 6 章回测不偷看未来的第一道闸。
时间轴(1 分钟 K 线示意) 09:00 ──▶ 09:01 ──▶ 09:02 ──▶ 09:03(现在) [闭合 ] [闭合 ] [闭合 ] [进行时 ✗ 不入正式表] └─▶ 到 09:04 才闭合落库
data_sources 模块(官方口径)承担连接管理、限速与重连这些公共职责;下面用一个骨架示意「接口数据到规范化字典」的旅程(教学示意,API 形态以官方文档为准):
# ws_consumer_demo.py —— 行情消费骨架示意(教学用,不可直接运行于实盘) import time def normalize(raw_msg): """把交易所原始消息映射为统一字典:字段名、单位、精度对齐""" return { "symbol": raw_msg.get("instId", "").upper(), # 统一符号写法 "ts": int(raw_msg.get("ts", 0)), # 毫秒时间戳,UTC "close": float(raw_msg.get("close", 0.0)), # 统一为浮点 "closed": bool(raw_msg.get("confirm", False)), # 闭合标志:闭合为准 } def on_message(raw_msg, sink): bar = normalize(raw_msg) if bar["closed"]: sink.append(bar) # 只有闭合根才进入落库通道 return bar if __name__ == "__main__": demo_feed = [ {"instId": "btc/usdt", "ts": 1767225600000, "close": "42000.5", "confirm": True}, {"instId": "btc/usdt", "ts": 1767225660000, "close": "42010.0", "confirm": False}, ] store = [] for msg in demo_feed: print(on_message(msg, store)) print("落库根数:", len(store)) # 期望:1(进行时根被挡下) time.sleep(0)
骨架里值得盯住的两行:normalize 负责把不同交易所的方言翻成统一普通话(符号、时间戳、精度——3.3 节的三重归一在此发端);closed 标志是闭合纪律的代码化身。真实模块还有限速器、重连状态机与订阅管理,读源码时按这四个概念对号入座即可。
把本节方法论压成一张可执行的走查清单——接入任何新数据源都按这五步走:
| 步 | 动作 | 完成判据 |
|---|---|---|
| 1 | 定通道:历史用 REST,实时用订阅 | 初始化与运行期各有明确来源 |
| 2 | 定数据:K 线为主,深度与成交按需 | 用途与存储预算写进设计说明 |
| 3 | 落四要点:限速、重连、心跳、闭合 | 每项都有配置或代码落点 |
| 4 | 试运行并人为制造一次断线 | 验证重连与补数路径真实可用 |
| 5 | 交质检:缺口检测与对账(3.3 节) | 质检报告没有未解释的缺口 |
第四步是清单里最容易被跳过、也最值钱的一步:断线补数路径平时不跑,等你发现它有毛病时,缺口往往已经攒了一屏。试运行阶段主动断一次网,成本最低、收益最大。
高频排查对照:
| 症状 | 常见原因 | 处置 |
|---|---|---|
| K 线数量与交易所页面不一致 | 进行时根被计入,或缺口未补 | 以闭合根为准;缺口走 REST 补洞 |
| 时间戳整体「差八小时」 | 本地时区与 UTC 混用 | 对内一律 UTC,展示层再转(3.3 节) |
| 频繁收到限速响应 | 客户端未节流或重试过猛 | 令牌桶加指数退避;429 当「停」信号 |
| 断线恢复后数据跳变 | 断线窗口未补齐 | 先补数再恢复消费,补不上进质检 |
行情接好了,但价格不是市场的全部——利率决议、新闻流、情绪指标同样驱动着市场。下一节看 data_providers 如何把这三类「慢变量」聚合成策略可用的输入。