Skip to content

17.2 实时 ETL 管道

名词解释:实时 ETL(Real-time ETL / Streaming ETL)

实时 ETL 是在数据产生的同时对其进行提取(Extract)、转换(Transform)和加载(Load)的过程。在量化交易中,实时 ETL 管道负责将来自交易所的原始 Tick 数据、公司公告、新闻等异源异构数据以低延迟转化为标准化、研究就绪的格式,供给策略引擎、风控系统和交易执行模块消费。

一、数据源接入:交易所直连、API、WebSocket

1.1 数据源的多样性

量化交易团队通常需要接入多种数据源:

数据源类型协议延迟典型数据
交易所直连FIX/OUCH/ITCH<10微秒订单簿、成交、委托确认
市场数据厂商TCP/UDP组播~100微秒聚合行情、NBBO
WebSocket APIWebSocket10-100毫秒加密货币行情、新闻
REST APIHTTP100-1000毫秒基本面、财报、另类数据
文件传输SFTP/S3分钟到小时日终批量数据

1.2 WebSocket 行情接入

PYTHON170 行 · 6.1 KB
📄此处有展示代码170 行 · 6.1 KB展开 ▼
python
import asyncio
import websockets
import json
import pandas as pd
from datetime import datetime
from typing import Callable, Dict, Any
import logging

class WebSocketMarketDataFeed:
    """
    基于 WebSocket 的实时行情数据接入。

    适用于:加密货币交易所、部分券商 API。
    """

    def __init__(self,
                 ws_url: str,
                 symbols: list,
                 on_trade: Callable = None,
                 on_quote: Callable = None,
                 on_error: Callable = None,
                 reconnect_delay: float = 1.0,
                 max_reconnect_attempts: int = 10):
        """
        参数:
            ws_url: WebSocket 端点 URL
            symbols: 订阅的标的列表
            on_trade: 成交回调函数(trade_data: dict)
            on_quote: 报价回调函数(quote_data: dict)
            on_error: 错误回调函数(error: Exception)
            reconnect_delay: 重连间隔(秒)
            max_reconnect_attempts: 最大重连次数
        """
        self.ws_url = ws_url
        self.symbols = symbols
        self.on_trade = on_trade or self._default_handler
        self.on_quote = on_quote or self._default_handler
        self.on_error = on_error or (lambda e: logging.error(f"WebSocket Error: {e}"))
        self.reconnect_delay = reconnect_delay
        self.max_reconnect_attempts = max_reconnect_attempts

        # 统计
        self.message_count = 0
        self.trade_count = 0
        self.quote_count = 0
        self.reconnect_count = 0

        # 各 symbol 的最新行情缓存
        self.latest_quotes: Dict[str, Dict[str, Any]] = {}
        self.latest_trades: Dict[str, Dict[str, Any]] = {}

    def _default_handler(self, data: dict):
        """默认处理:存储到缓存"""
        if data.get('type') == 'trade':
            self.latest_trades[data['symbol']] = data
            self.trade_count += 1
        elif data.get('type') == 'quote':
            self.latest_quotes[data['symbol']] = data
            self.quote_count += 1

        self.message_count += 1

    async def _subscribe(self, websocket):
        """发送订阅消息(根据具体交易所 API 实现)"""
        subscribe_msg = {
            "method": "SUBSCRIBE",
            "params": [f"{s.lower()}@trade" for s in self.symbols] +
                      [f"{s.lower()}@depth" for s in self.symbols],
            "id": 1
        }
        await websocket.send(json.dumps(subscribe_msg))
        logging.info(f"Subscribed to {len(self.symbols)} symbols")

    async def _message_parser(self, raw_message: str) -> dict:
        """
        解析原始 WebSocket 消息为标准化格式。

        此为简化示例,实际需根据具体数据源的格式适配。
        """
        try:
            msg = json.loads(raw_message)

            # 示例:将不同交易所的格式统一为标准格式
            if 'e' in msg and msg['e'] == 'trade':
                # Binance 成交格式
                return {
                    'type': 'trade',
                    'symbol': msg['s'],
                    'price': float(msg['p']),
                    'size': float(msg['q']),
                    'timestamp': datetime.fromtimestamp(msg['T'] / 1000),
                    'trade_id': msg['t'],
                    'is_buyer_maker': msg['m']
                }
            elif 'type' in msg and msg['type'] == 'ticker':
                # Coinbase 格式
                return {
                    'type': 'quote',
                    'symbol': msg['product_id'],
                    'bid': float(msg['best_bid']),
                    'ask': float(msg['best_ask']),
                    'bid_size': float(msg.get('best_bid_size', 0)),
                    'ask_size': float(msg.get('best_ask_size', 0)),
                    'timestamp': datetime.fromisoformat(
                        msg['time'].replace('Z', '+00:00')
                    )
                }
            else:
                return msg

        except json.JSONDecodeError:
            return {'type': 'unknown', 'raw': raw_message}
        except Exception as e:
            logging.error(f"Message parse error: {e}")
            return {'type': 'error', 'error': str(e)}

    async def _process_messages(self, websocket):
        """处理 WebSocket 消息流"""
        async for raw_message in websocket:
            try:
                parsed = await self._message_parser(raw_message)

                if parsed.get('type') == 'trade':
                    self.on_trade(parsed)
                elif parsed.get('type') == 'quote':
                    self.on_quote(parsed)

            except Exception as e:
                self.on_error(e)

    async def connect(self):
        """建立 WebSocket 连接并开始消费数据"""
        attempt = 0

        while attempt < self.max_reconnect_attempts:
            try:
                async with websockets.connect(
                    self.ws_url,
                    ping_interval=20,
                    ping_timeout=10,
                    max_size=2**24  # 16MB max message size
                ) as ws:
                    logging.info(f"Connected to {self.ws_url}")
                    attempt = 0  # 重置重连计数

                    await self._subscribe(ws)
                    await self._process_messages(ws)

            except websockets.ConnectionClosed as e:
                self.reconnect_count += 1
                attempt += 1
                logging.warning(f"Connection closed: {e}. "
                                f"Reconnect attempt {attempt}/{self.max_reconnect_attempts}")
                await asyncio.sleep(self.reconnect_delay * attempt)

            except Exception as e:
                self.on_error(e)
                attempt += 1
                await asyncio.sleep(self.reconnect_delay)

        logging.error("Max reconnect attempts reached. Giving up.")

    def get_statistics(self) -> dict:
        """获取数据统计"""
        return {
            'message_count': self.message_count,
            'trade_count': self.trade_count,
            'quote_count': self.quote_count,
            'reconnect_count': self.reconnect_count
        }
点击展开可浏览运行结果
📘 本段代码定义了 1 个函数/类:类 `WebSocketMarketDataFeed`(基于 WebSocket 的实时行情数据接入)。该片段为教学展示(未包含独立运行的输入数据),可在实战练习中结合真实数据调用。

2.1 流处理架构

典型的量化实时 ETL 架构:

[交易所/数据源] --> [Kafka Topic] --> [Flink/Kafka Streams]
                                          |
                               +----------+-----------+
                               |                      |
                          [实时特征计算]         [异常检测&告警]
                               |                      |
                          [Redis/特征存储]      [监控面板/企业微信]
                               |
                          [策略引擎消费]

2.2 Kafka Streams 数据规范化

PYTHON87 行 · 2.9 KB
📄此处有展示代码87 行 · 2.9 KB展开 ▼
python
from kafka import KafkaProducer, KafkaConsumer
import msgpack
import time

class KafkaTickPipeline:
    """
    基于 Kafka 的 Tick 数据流处理管道。

    Topic 设计:
    - raw.ticks.{exchange}      : 原始数据(直接接入)
    - normalized.ticks.{asset}   : 标准化数据(清洗后)
    - derived.bars.{interval}    : 聚合 K 线
    - alerts.anomaly             : 异常告警
    """

    def __init__(self, bootstrap_servers: list):
        self.bootstrap_servers = bootstrap_servers

        # Producer: 序列化用 msgpack(比 JSON 快 5-10x)
        self.producer = KafkaProducer(
            bootstrap_servers=bootstrap_servers,
            value_serializer=lambda v: msgpack.packb(v, use_bin_type=True),
            compression_type='lz4',
            linger_ms=5,          # 批处理延迟
            batch_size=16384      # 批大小
        )

    def normalize_tick(self, raw_tick: dict, source: str) -> dict:
        """
        将不同来源的 Tick 数据标准化为统一 Schema。

        参数:
            raw_tick: 原始 tick 字典
            source: 数据来源标识(如 'binance', 'xtp', 'ctp')
        返回:
            标准化的 tick 字典
        """
        # 统一 Schema 定义
        normalized = {
            'timestamp': None,
            'symbol': None,
            'exchange': None,
            'price': None,
            'size': None,
            'bid': None,
            'ask': None,
            'bid_size': None,
            'ask_size': None,
            'trade_condition': None,
            'source': source,
            'arrival_time': int(time.time() * 1e9)  # 本系统接收时间
        }

        # 按来源适配
        if source == 'binance':
            normalized['timestamp'] = raw_tick.get('T', 0)
            normalized['symbol'] = raw_tick.get('s', '')
            normalized['price'] = float(raw_tick.get('p', 0))
            normalized['size'] = float(raw_tick.get('q', 0))

        elif source == 'xtp':
            normalized['timestamp'] = raw_tick.get('data_time', 0)
            normalized['symbol'] = raw_tick.get('ticker', '')
            normalized['price'] = raw_tick.get('last_price', 0)
            normalized['size'] = raw_tick.get('qty', 0)
            normalized['bid'] = raw_tick.get('bid', [None])[0]
            normalized['ask'] = raw_tick.get('ask', [None])[0]

        # 数据质量检查
        if normalized['price'] is not None and normalized['price'] <= 0:
            return None  # 过滤无效价格

        return normalized

    def publish_normalized_tick(self, tick: dict, asset_class: str):
        """发布标准化 Tick 到 Kafka"""
        if tick is None:
            return

        topic = f"normalized.ticks.{asset_class}"
        key = tick['symbol'].encode('utf-8')

        self.producer.send(topic, key=key, value=tick)

    def close(self):
        self.producer.flush()
        self.producer.close()
点击展开可浏览运行结果
📘 本段代码定义了 1 个函数/类:类 `KafkaTickPipeline`(基于 Kafka 的 Tick 数据流处理管道)。该片段为教学展示(未包含独立运行的输入数据),可在实战练习中结合真实数据调用。

三、数据规范化与时间对齐

3.1 时间对齐的挑战

在跨市场和跨资产数据中,时间对齐是一个核心挑战:

  • 交易所时钟差异:不同交易所的系统时钟可能偏差几十毫秒到几秒
  • NTP 同步:即使使用 NTP,精度也限于毫秒级别
  • 数据抵达顺序:网络延迟导致数据可能乱序抵达(out-of-order arrival)
PYTHON36 行 · 1.2 KB
📄此处有展示代码36 行 · 1.2 KB展开 ▼
python
def timestamp_alignment(events: pd.DataFrame,
                         tolerance_ms: int = 100) -> pd.DataFrame:
    """
    对来自不同数据源的异源事件进行时间对齐。

    策略:使用滑动窗口,将容差内的事件视为"同时发生"。

    参数:
        events: 多源事件DataFrame(含 source, timestamp, symbol 列)
        tolerance_ms: 容差(毫秒)
    返回:
        对齐后的事件DataFrame(新增 aligned_ts 列)
    """
    tolerance_ns = pd.Timedelta(milliseconds=tolerance_ms)

    # 对每个 symbol,按时间排序
    aligned_frames = []

    for symbol, group in events.groupby('symbol'):
        group = group.sort_values('timestamp')

        # 将容差内的事件对齐到最早时间
        group['time_diff'] = group['timestamp'].diff()

        # 创建对齐组:当时间差超过容差时,开始新的对齐组
        group['aligned_group'] = (
            group['time_diff'] > tolerance_ns
        ).cumsum()

        # 每个对齐组的时间基准为组内最早时间
        group['aligned_ts'] = group.groupby('aligned_group')['timestamp'] \
            .transform('min')

        aligned_frames.append(group)

    return pd.concat(aligned_frames)
点击展开可浏览运行结果
📘 本段代码定义了 1 个函数/类:函数 `timestamp_alignment`(对来自不同数据源的异源事件进行时间对齐)。该片段为教学展示(未包含独立运行的输入数据),可在实战练习中结合真实数据调用。

四、延迟监控与容错

4.1 端到端延迟度量

PYTHON52 行 · 2.1 KB
📄此处有展示代码52 行 · 2.1 KB展开 ▼
python
class LatencyMonitor:
    """
    实时 ETL 管道的延迟监控。

    度量维度:
    1. 接入延迟:交易所生成时间 -> 本系统接收时间
    2. 处理延迟:接收时间 -> 标准化完成时间
    3. 发布延迟:标准化完成 -> 下游消费时间
    4. 端到端延迟:交易所生成 -> 策略引擎消费
    """

    def __init__(self, window_seconds: int = 60):
        self.window_seconds = window_seconds
        self.ingest_latencies = []    # 接入延迟
        self.process_latencies = []   # 处理延迟
        self.publish_latencies = []   # 发布延迟
        self.e2e_latencies = []       # 端到端延迟
        self.message_counts = []      # 消息量

    def record(self,
               exchange_ts: int,      # 交易所时间戳(纳秒)
               arrival_ts: int,       # 本系统接收时间(纳秒)
               processed_ts: int,     # 处理完成时间(纳秒)
               consumed_ts: int = None):  # 下游消费时间(纳秒)
        """记录一条消息的延迟"""
        ingest_lat = (arrival_ts - exchange_ts) / 1e6  # 转为毫秒
        process_lat = (processed_ts - arrival_ts) / 1e6

        self.ingest_latencies.append(ingest_lat)
        self.process_latencies.append(process_lat)

        if consumed_ts is not None:
            pub_lat = (consumed_ts - processed_ts) / 1e6
            e2e_lat = (consumed_ts - exchange_ts) / 1e6
            self.publish_latencies.append(pub_lat)
            self.e2e_latencies.append(e2e_lat)

    def get_percentiles(self) -> dict:
        """获取延迟百分位数"""
        def safe_percentile(data, p):
            if not data:
                return None
            return np.percentile(data, p)

        return {
            'ingest_p50': safe_percentile(self.ingest_latencies, 50),
            'ingest_p99': safe_percentile(self.ingest_latencies, 99),
            'process_p50': safe_percentile(self.process_latencies, 50),
            'process_p99': safe_percentile(self.process_latencies, 99),
            'e2e_p50': safe_percentile(self.e2e_latencies, 50),
            'e2e_p99': safe_percentile(self.e2e_latencies, 99),
        }
点击展开可浏览运行结果
📘 本段代码定义了 1 个函数/类:类 `LatencyMonitor`(实时 ETL 管道的延迟监控)。该片段为教学展示(未包含独立运行的输入数据),可在实战练习中结合真实数据调用。

4.2 容错策略

实时 ETL 管道的容错设计采用多层防护:

  1. 数据校验层:检查字段完整性、价格合理性、时间单调性
  2. 死信队列(DLQ):无法处理的消息进入 DLQ,人工或自动化修复
  3. 断点续传:Kafka 的 offset 管理保证崩溃恢复后不会丢失或重复数据
  4. 多路冗余:关键数据源使用多个厂商交叉验证

实时 ETL 的最终目标是确保策略引擎在任何时刻都能获得高质量、低延迟的数据。而数据质量的持续保障——包括完整性检查、异常检测和版本管理——是我们接下来要讨论的数据质量监控体系的核心课题。