主题切换
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 API | WebSocket | 10-100毫秒 | 加密货币行情、新闻 |
| REST API | HTTP | 100-1000毫秒 | 基本面、财报、另类数据 |
| 文件传输 | SFTP/S3 | 分钟到小时 | 日终批量数据 |
1.2 WebSocket 行情接入
此处有展示代码展开 ▼
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 的实时行情数据接入)。该片段为教学展示(未包含独立运行的输入数据),可在实战练习中结合真实数据调用。
二、流处理框架(Flink/Kafka Streams)
2.1 流处理架构
典型的量化实时 ETL 架构:
[交易所/数据源] --> [Kafka Topic] --> [Flink/Kafka Streams]
|
+----------+-----------+
| |
[实时特征计算] [异常检测&告警]
| |
[Redis/特征存储] [监控面板/企业微信]
|
[策略引擎消费]2.2 Kafka Streams 数据规范化
此处有展示代码展开 ▼
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)
此处有展示代码展开 ▼
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 端到端延迟度量
此处有展示代码展开 ▼
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 管道的容错设计采用多层防护:
- 数据校验层:检查字段完整性、价格合理性、时间单调性
- 死信队列(DLQ):无法处理的消息进入 DLQ,人工或自动化修复
- 断点续传:Kafka 的 offset 管理保证崩溃恢复后不会丢失或重复数据
- 多路冗余:关键数据源使用多个厂商交叉验证
实时 ETL 的最终目标是确保策略引擎在任何时刻都能获得高质量、低延迟的数据。而数据质量的持续保障——包括完整性检查、异常检测和版本管理——是我们接下来要讨论的数据质量监控体系的核心课题。