行情数据接入与清洗

行情数据接入与清洗的完整工程实践:Level1 与 Level2 数据的差异、CTP 与交易所协议解析、快照与增量订单簿重建、时间戳与延迟测量、异常 tick 清洗规则、除权除息复权、Parquet 列式存储与压缩、历史回放以及多源校验与容灾方案。

行情是量化系统的地基。地基不牢,上层的因子、回测、执行全部失真。工程上真正困难的地方不是「怎么连上行情源」——CTP 的接入文档写得很清楚——而是如何保证到达策略手里的每一条数据都是干净、及时、可复现的。脏数据的破坏力是静默的:它不会报错,只会让你的回测结果悄悄偏离实盘。

行情数据接入有三个层次的问题需要分开处理:传输层(协议解析、丢包重传)、语义层(快照与增量的含义、时间戳的定义)、存储层(格式、压缩、回放一致性)。很多团队把三者混在一个脚本里,结果行情一抖动就全链路崩溃。

本文按数据流向展开:先讲数据层次与协议,再讲订单簿重建与时间戳,然后是清洗规则、复权、存储、回放,最后是多源校验与容灾。这是 量化交易系统全景 里数据层的第一篇下钻。

目录

  1. 行情数据的三个层次
  2. 交易所协议与接入方式
  3. 快照与增量:订单簿重建
  4. 时间戳与延迟测量
  5. 异常 tick 清洗规则
  6. 复权与除权除息
  7. 存储格式与压缩
  8. 历史回放与仿真
  9. 多源校验与容灾

1. 行情数据的三个层次

不同市场提供的行情粒度差异极大,理解这三个层次是选型的前提:

层次内容频率典型来源数据量/日
Level1最优买卖价 + 成交量3 秒快照免费行情~5 GB(全市场)
Level2十档盘口 + 逐笔委托/成交实时推送付费行情~50 GB
逐笔委托每一笔报单/撤单实时交易所直发~100 GB

Level1 只有一档,适合中低频策略;Level2 有十档盘口,是做市和日内策略的最低要求;逐笔委托(order-by-order)能重建完整订单簿,是高频策略的必需品,但价格昂贵且需要交易所级别的接入权限。

选择依据是策略需要看到多少流动性。一个只用到收盘价的策略用 Level1 就够;一个需要判断「盘口压力」的策略必须看多档;一个需要预测「下一笔成交方向」的策略必须看逐笔。

2. 交易所协议与接入方式

国内主流接入方式:

协议/接口市场特点延迟
CTP商品/金融期货事实标准,C++ API毫秒级
券商极速柜台股票各券商自研亚毫秒
交易所 Binary股票 L2二进制,需报备微秒级
FAST国际模板化压缩微秒级
ITCH/OUCH美股逐笔 + 报单微秒级

CTP 的行情回调是全量快照:

// CTP 行情回调:每个 tick 都是完整快照
void OnRtnDepthMarketData(CThostFtdcDepthMarketDataField* p) {
    Tick t;
    t.symbol     = p->InstrumentID;
    t.last_price = p->LastPrice;
    t.bid1       = p->BidPrice1;
    t.ask1       = p->AskPrice1;
    t.bid_vol1   = p->BidVolume1;
    t.ask_vol1   = p->AskVolume1;
    t.volume     = p->Volume;          // 累计成交量,非增量
    t.turnover   = p->Turnover;        // 累计成交额
    t.update_ms  = p->UpdateMillisec;
    t.action_day = p->ActionDay;       // 交易日,夜盘跨日关键
    dispatcher_.push(std::move(t));    // 推入无锁队列
}

注意 Volume 是累计值而不是增量,必须自己差分。夜盘品种的 ActionDay 与自然日不同,日期归属必须用 ActionDay。

3. 快照与增量:订单簿重建

Level2 行情通常是增量的:交易所只推送变化的档位或逐笔事件,客户端负责维护本地订单簿。重建逻辑的正确性直接决定策略能否看到真实盘口。

class OrderBook:
    def __init__(self):
        self.bids = {}   # price -> qty,买单
        self.asks = {}   # price -> qty,卖单
        self.last_seq = 0

    def apply(self, msg):
        if msg.seq != self.last_seq + 1:   # 序号校验:丢包立即触发快照重建
            raise SequenceGap(self.last_seq, msg.seq)
        if msg.type == 'ADD':
            book = self.bids if msg.side == 'B' else self.asks
            book[msg.price] = book.get(msg.price, 0) + msg.qty
        elif msg.type == 'CANCEL':
            book = self.bids if msg.side == 'B' else self.asks
            book[msg.price] = max(0, book.get(msg.price, 0) - msg.qty)
            if book[msg.price] == 0:
                del book[msg.price]
        elif msg.type == 'TRADE':
            book = self.asks if msg.side == 'B' else self.bids
            book[msg.price] = max(0, book.get(msg.price, 0) - msg.qty)
        self.last_seq = msg.seq

    def best_bid(self):
        return max(self.bids) if self.bids else None

    def best_ask(self):
        return min(self.asks) if self.asks else None

关键点:序号连续性校验是唯一的丢包检测手段。一旦发现序号跳跃,必须立即向交易所请求快照重建,而不是继续用错误的订单簿交易。

4. 时间戳与延迟测量

一条行情至少有三个时间:

时间含义用途
交易所时间事件发生时刻策略逻辑的时间基准
网关接收时间你的机器收到的时刻测量传输延迟
策略处理时间策略看到数据的时刻测量本地处理延迟

延迟 = 策略处理时间 − 交易所时间。没有精确的本地时钟(PTP 或硬件时间戳),这个差值毫无意义。

// 用硬件时间戳网卡测量行情延迟
struct TimedTick {
    Tick tick;
    uint64_t hw_rx_ns;   // 网卡硬件打戳,纳秒
    uint64_t sw_rx_ns;   // 用户态收到时刻
    uint64_t done_ns;    // 策略处理完成时刻
};

void on_packet(const uint8_t* buf, size_t len, uint64_t hw_ns) {
    TimedTick t;
    t.hw_rx_ns = hw_ns;
    t.sw_rx_ns = now_ns();
    parse(buf, len, t.tick);
    t.done_ns = now_ns();
    metrics_.record("md.total_ns", t.done_ns - t.hw_rx_ns);
}

监控这条延迟的 P99 比平均值重要得多,因为策略的触发往往是延迟尖峰导致的。

5. 异常 tick 清洗规则

真实行情里混着大量「看起来正常」的脏数据。必须成体系地清洗:

异常类型现象处理
价格越界超出涨跌停丢弃或截断
买卖价倒挂bid > ask丢弃
零成交量volume 不变但 last 变保留但标记
时间倒流ts 小于上一条按序重排
重复推送与上一条完全相同去重
涨跌停虚假挂单封板时巨量挂单识别并降权
def clean_tick(t, prev, limits):
    if t.last_price <= 0 or t.last_price > limits.upper:
        return None                       # 价格越界
    if t.bid1 > 0 and t.ask1 > 0 and t.bid1 >= t.ask1:
        return None                       # 买卖倒挂
    if prev and t.ts <= prev.ts:
        t.ts = prev.ts + 1                # 时间倒流,单调化
    if prev and t.volume < prev.volume:
        return None                       # 累计量回退,疑似重启
    return t

涨跌停虚假挂单是最难处理的:封板时盘口会有几十万手的挂单,这些单子大概率不会成交,但如果回测按盘口量估算可成交量,会严重高估。常见做法是识别「价格 = 涨跌停价」的档位并打折或剔除。

6. 复权与除权除息

股票行情必须处理除权除息,否则会出现「一夜之间跌 30%」的假信号。三种复权方式:

方式做法适用
前复权以最新价为基准调整历史技术分析、回测
后复权以最早价为基准调整未来长期收益计算
不复权保持原始价撮合、成交价

回测用前复权最自然,因为最新价不变、历史价被缩放。但要注意:前复权价格会随新的除权事件变化,导致回测结果不可复现——今天跑的历史数据和昨天跑的不一样。严肃的回测应该用后复权,因为后复权价一旦确定就不再变化。

def adjust_backward(bars, dividends):
    factor = 1.0   # 后复权:从最早一天开始,逐日累乘复权因子
    for bar in bars:
        if bar.ex_date in dividends:
            d = dividends[bar.ex_date]
            factor *= (bar.close + d.cash) / bar.close
        bar.adj_close = bar.close * factor
    return bars

期货不存在复权问题,但存在换月(主力合约切换),需要在切换点做价格拼接,否则跨月回测会出现假跳空。

7. 存储格式与压缩

tick 数据不适合行式数据库。列式存储的优势:同列数据类型一致、压缩率高、只读需要的列。

格式压缩比随机读生态
CSV1x差通用
Parquet5~10x中Python/Spark
Arrow IPC3~5x极好内存映射
ClickHouse8~15x极好SQL 查询

推荐分层:热数据用 Arrow IPC 做内存映射(回放时零拷贝),温数据用 Parquet,冷数据压缩后归档。

import pyarrow as pa
import pyarrow.parquet as pq

table = pa.Table.from_pylist(ticks, schema=TICK_SCHEMA)
pq.write_table(
    table,
    f"ticks/{symbol}/{date}.parquet",
    compression="zstd",
    compression_level=3,      # 3 是速度/压缩比的最佳平衡点
    row_group_size=100_000,   # 每 10 万行一个 row group
)

row_group_size 决定随机读的粒度,太小则元数据开销大,太大则单次读取浪费 IO。10 万行是个常用起点。

8. 历史回放与仿真

回放的目标是让回测消费与实盘完全一致的数据流。两种模式:

加速回放:把历史 tick 按压缩时间轴快速灌入(用于回测)
实时回放:按原始时间间隔灌入(用于仿真盘)
class ReplayFeed:
    def __init__(self, path, speed=0):
        self.ticks = pq.read_table(path).to_pylist()
        self.speed = speed        # 0 表示全速,1 表示实时

    def run(self, on_tick):
        base_ts = self.ticks[0]['ts']
        t0 = time.time_ns()
        for t in self.ticks:
            if self.speed > 0:
                target = t0 + int((t['ts'] - base_ts) / self.speed)
                while time.time_ns() < target:
                    pass          # 自旋等待,保证时序
            on_tick(t)

回放的一致性要求:同一份数据文件,回测与仿真盘产生相同的 tick 序列。做不到这一点,仿真盘就失去了验证价值。

9. 多源校验与容灾

单一行情源是不可接受的单点故障。生产系统至少要有主备两路:

方案切换时间复杂度适用
主备冷切换秒级低中低频
双路热备毫秒级中中高频
三路投票微秒级高做市

双路热备要解决去重问题:两路行情到达时间不同,需要按「交易所序号 + 时间戳」去重。

class DualFeedMerger:
    def __init__(self):
        self.seen = LRUSet(maxsize=100_000)
        self.primary_ok = True

    def merge(self, t, source):
        key = (t.symbol, t.ts, t.last_price, t.volume)
        if key in self.seen:
            return None                 # 重复,丢弃
        self.seen.add(key)
        return t

行情源通过 Kafka 汇聚时,多路数据的顺序保证可以参考 Kafka 消费者组再均衡 里的分区策略——按 symbol 分区才能保证单标的时序。整体架构上,行情管道的设计原则与 Kafka 入门 中描述的消息总线模式高度一致。

权衡取舍

维度Level1Level2逐笔
成本低高极高
延迟3 秒实时实时
策略上限中低频日内高频
存储压力小大极大

清洗策略上也有取舍:激进清洗(丢弃一切可疑数据)能保证干净但会丢失真实的市场异动;保守清洗(只标记不丢弃)保留信息但把判断责任推给策略。生产系统通常采用「标记 + 分层」:核心字段严格校验,辅助字段宽松处理并打上质量标记。

复权方式的选择上,前复权直观但不可复现,后复权可复现但价格不直观。回测与因子研究用后复权,展示与图表用前复权,是较稳妥的组合。

常见坑清单

  1. 把累计成交量当增量用:CTP 的 Volume 是累计值,直接当增量会算出天量成交量。
  2. 忽略夜盘日期归属:用自然日而非 ActionDay,夜盘数据会归到错误的一天。
  3. 订单簿不做序号校验:丢包后订单簿永久失真,策略基于错误盘口下单。
  4. 前复权做回测:新的除权事件会改变历史数据,回测结果不可复现。
  5. 时间戳用本地时间:未做时区与交易日处理,跨市场数据无法对齐。
  6. 单一行情源无备份:行情源一断,策略直接停摆或基于过期数据下单。
  7. CSV 存 tick:单日全市场数据上百 GB,CSV 读取慢到无法回测。
  8. 清洗规则硬编码:涨跌停价随市场变化,硬编码会在极端行情下失效。
  9. 回放不保证时序:全速回放让策略看到「同一纳秒」的多条数据,与实盘不符。
  10. 不做数据质量监控:脏数据静默污染下游,直到实盘亏损才发现。

小结

行情数据接入的核心是三条保证:数据干净(清洗规则成体系)、时序正确(时间戳与序号校验)、可复现(回放与实盘同源)。做到这三条,上层的回测才有意义。

工程上最容易被低估的是「清洗」这一环。它不像低延迟那样有戏剧性,但决定了策略的真实性。一个把涨跌停虚假挂单算进可成交量的回测,年化可以虚高一倍以上。

下一步建议阅读 回测框架设计与前视偏差 ,看干净的数据如何被正确地使用。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「量化交易」更多文章

  1. 风险模型与因子归因
  2. 回测偏差与过拟合防范
  3. 市场微结构与流动性