通过内容哈希+单调序列号实现WebSocket市场数据采集幂等性,避免网络抖动导致重复数据污染回测系统。
大多数自建交易数据管道在演示的成功路径下看起来都是正确的,但一旦网络出现抖动,就会悄悄破坏回测结果。罪魁祸首几乎总是一样的:WebSocket 重连时,feed 会重新传输客户端已经收到的 ticks,而如果没有去重层,这些重复数据就会直接流入你的特征存储和标签。
本文展示了一个小型、可复现的模式——内容哈希加上每流单调递增的序列号——这使得收集器具有幂等性:处理同一条 tick 两次产生与处理一次相同的状态。
WebSocket 流不是具有交付保证的队列。当 socket 在 tick 批次中间断开并重连时,常见行为有:
这些在 30 秒的本地测试中都不可见。它们在凌晨 3 点浮现——当交易所滚动序列号时。
对于每条标准化的 tick,从定义 tick 标识的字段(instrument、timestamp-or-trade-id、price、size、side)计算一个稳定的内容哈希。保持一个每流 last_seq 单调递增计数器。在 ingest 时:
如果 tick.seq <= last_seq[stream] → 丢弃(已见过 / 乱序)。
否则计算 h = hash(normalized(tick))。如果 h 在最近见过集合中(有界 LRU)→ 丢弃。
否则接受,更新 last_seq 和见过集合,并持久化。
内容哈希可以捕获语义上完全相同的重复,即使序列号不可靠;序列号可以捕获乱序重传。两者都需要。
最小参考实现(Python,provider 无关)
import hashlib
from collections import defaultdict, deque
class IdempotentCollector:
def __init__(self, seen_window: int = 100_000):
self.last_seq = defaultdict(int)
self.seen = defaultdict(lambda: deque(maxlen=seen_window))
def _identity_hash(self, tick) -> str:
key = f"{tick['sym']}|{tick['ts']}|{tick['px']}|{tick['sz']}|{tick['side']}"
return hashlib.sha256(key.encode()).hexdigest()[:16]
def ingest(self, stream: str, tick: dict) -> dict | None:
seq = tick.get("seq", 0)
if seq and seq <= self.last_seq[stream]:
return None # reorder / duplicate by sequence
h = self._identity_hash(tick)
if h in self.seen[stream]:
return None # duplicate by content
self.last_seq[stream] = max(self.last_seq[stream], seq)
self.seen[stream].append(h)
return tick # accepted exactly once
故障注入测试框架
要证明幂等性,必须注入故障,而不是希望它们不发生:
Timeout:在批次中间断开 socket,重连,重放尾部窗口。断言 accepted-count == unique-count。
Duplicate POST:在同一批次中发送同一条 tick 两次。断言一次接受。
Reconnect reorder:先交付 ticks 5,6,7,然后再交付 4,5,6。断言 4 被接受一次,5/6/7 各一次。
运行 N 个周期;如果每次运行中 accepted_total == distinct_ticks,则收集器通过。这是手册中幂等性规则(§11.2)在任何非盲重试之前要求的测试。
seen-window 是有界的;极长寿的流理论上可能让超出窗口的重复通过。根据你的 feed 的重传窗口来设置其大小。
Provider 特定行为(snapshot vs delta、resume token vs sequence)仍然需要每个 provider 的说明。本模式假设一个有序或可重传的 feed。
这是一种工程控制,不是交易信号。它使你的数据值得信赖;它没有说明在其上训练的模型是否赚钱。
如果你想要完整的测试框架 + 结果表,请参阅配套的 GitHub 仓库(在发布时添加链接)。对于更广泛的数据管道,请阅读关于 walk-forward / out-of-sample 原则的回测笔记。