第二十四章:系统架构设计:实时数据管道、计算引擎、信号分发系统
做波动率曲面交易,最怕什么?
不是策略亏钱,是信号出来了,系统没跟上。
我见过太多团队,策略逻辑写得漂漂亮亮,一到实盘就崩。数据延迟、计算超时、信号发不出去……说白了,架构没撑住。今天我们就聊聊,一个能打的波动率曲面交易系统,到底该怎么搭。
24.1 整体架构:三驾马车
我个人习惯把系统拆成三个核心模块:
- 实时数据管道:负责把行情数据从交易所拉到本地,清洗、对齐、缓存
- 计算引擎:负责曲面构建、特征提取、信号生成
- 信号分发系统:负责把交易信号送到执行端,同时做好风控
这三个模块,缺一个都不行。你想想看,数据管道断了,计算引擎就是无米之炊;计算引擎慢了,信号出来行情都走完了;分发系统不稳,信号到了执行端也白搭。
核心原则:每个模块都要能独立部署、独立扩容、独立降级。我在项目中遇到过,某次数据源挂了,幸好管道层有缓存,计算引擎还能用历史数据跑模拟信号,没让整个系统停摆。
24.2 实时数据管道:别让数据成为瓶颈
数据管道这块,我踩过的坑最多。刚开始做的时候,觉得不就是收个行情嘛,有什么难的?结果一上实盘,各种问题全冒出来了。
24.2.1 数据接入层
期权数据比股票复杂得多。每个到期日、每个行权价,都是独立的合约。一个活跃的期权品种,光快照数据每秒就能产生几千条。
我个人习惯用这样的分层设计:
# 伪代码:数据管道接入层
class DataIngestionLayer:
def __init__(self):
self.raw_queue = asyncio.Queue(maxsize=10000)
self.clean_queue = asyncio.Queue(maxsize=5000)
async def ingest(self, source):
# 从交易所接收原始数据
async for msg in source.stream():
await self.raw_queue.put(msg)
async def clean(self):
while True:
raw = await self.raw_queue.get()
# 去重、时间戳对齐、异常值过滤
cleaned = self._cleanse(raw)
await self.clean_queue.put(cleaned)
经验之谈:我曾经因为没做去重,导致同一个tick被计算了两次,曲面瞬间扭曲,系统直接发了假信号。从那以后,我每条数据都加上了sequence number,重复的一律丢弃。
24.2.2 数据对齐与缓存
不同合约的数据到达时间不一样。近月合约流动性好,数据来得快;远月合约可能几秒才更新一次。直接拿最新数据拼曲面,会出问题。
我的做法是:维护一个时间窗口(比如100毫秒),窗口内的数据对齐到同一个时间戳,再送进计算引擎。
| 组件 | 技术选型 | 说明 |
|---|---|---|
| 消息队列 | Redis Stream / Kafka | 缓冲数据,防止计算引擎被冲垮 |
| 内存数据库 | Redis | 缓存最新快照,支持快速查询 |
| 时序数据库 | InfluxDB / ClickHouse | 存储历史数据,用于回测和分析 |
24.3 计算引擎:曲面构建与信号生成
计算引擎是整个系统的心脏。它要完成三件事:构建曲面、提取特征、生成信号。
24.3.1 曲面构建模块
收到对齐后的数据,第一步就是构建波动率曲面。我一般用SVI模型或者样条插值。这里要注意,每次收到新数据,不是全量重建,而是增量更新。
# 伪代码:增量曲面更新
class VolSurfaceEngine:
def __init__(self):
self.surface = None # 当前曲面
self.last_update = 0
def update(self, new_quotes):
# 只更新有变化的节点
changed_nodes = self._detect_changes(new_quotes)
if not changed_nodes:
return # 没变化,跳过
# 局部插值,避免全量计算
self.surface = self._local_interpolate(changed_nodes)
self.last_update = time.time()
# 检查曲面是否合理
if not self._validate_surface():
self._fallback_to_previous() # 曲面异常,回退
注意:曲面构建一定要做合理性校验。我遇到过,某次因为数据源推送了错误的价格,曲面瞬间出现负的隐含波动率。幸好校验层拦住了,没让信号发出去。
24.3.2 特征提取与信号生成
曲面构建好了,接下来就是提取特征。常用的特征包括:
- 曲面斜率(skew)
- 曲面曲率(convexity)
- 期限结构(term structure)
- 局部波动率 vs 隐含波动率的偏离
信号生成这块,我习惯用规则引擎+轻量模型结合的方式。规则引擎处理明显的套利机会,模型处理更复杂的模式识别。
# 伪代码:信号生成
class SignalGenerator:
def __init__(self):
self.rules = [SkewRule(), TermStructureRule(), ArbitrageRule()]
self.ml_model = load_model('surface_pattern_model.pkl')
def generate(self, surface):
signals = []
# 规则引擎
for rule in self.rules:
sig = rule.evaluate(surface)
if sig:
signals.append(sig)
# 模型预测
features = self._extract_features(surface)
ml_signal = self.ml_model.predict(features)
if ml_signal.confidence > 0.7:
signals.append(ml_signal)
return signals
24.4 信号分发系统:最后一公里
信号生成后,怎么送到执行端?这里面的门道不少。
24.4.1 信号路由
不同的信号要去不同的地方。套利信号可能直接送OMS,风控信号要送监控系统,统计信号可能先送人工审核队列。
我一般用发布-订阅模式:
- 每个信号带一个topic标签
- 订阅者根据topic决定是否处理
- 支持多级路由,比如先过风控,再过执行
24.4.2 风控检查
信号发出前,必须过风控。这是最后一道防线。
| 风控规则 | 检查内容 | 触发动作 |
|---|---|---|
| 价格合理性 | 信号价格是否偏离市场太大 | 拦截并告警 |
| 仓位限制 | 当前仓位是否超过阈值 | 拒绝或缩减规模 |
| 频率限制 | 同一信号是否频繁触发 | 冷却期 |
| 相关性检查 | 是否与其他持仓高度相关 | 人工确认 |
一个小技巧:风控检查最好做成可配置的。不同策略、不同市场环境,风控参数不一样。我习惯把规则写在YAML文件里,系统启动时加载,运行时也能热更新。
24.4.3 执行对接
信号通过风控后,就要送到执行系统。这里要注意协议兼容性。有的交易所用FIX协议,有的用REST API,还有的用WebSocket。
我的做法是:在分发系统里做一个适配层,统一信号格式,底层对接不同的执行通道。
# 伪代码:执行适配层
class ExecutionAdapter:
def __init__(self):
self.adapters = {
'fix': FIXAdapter(),
'rest': RESTAdapter(),
'ws': WebSocketAdapter()
}
def execute(self, signal):
adapter = self.adapters[signal.exchange_type]
result = adapter.send(signal)
return result
24.5 性能与可靠性
最后聊两句性能。波动率曲面交易对延迟要求没那么变态,但也不能太慢。我一般要求:
- 数据管道:端到端延迟 < 50ms
- 曲面计算:单次更新 < 10ms
- 信号分发:从生成到送达 < 5ms
可靠性方面,每个模块都要有降级方案。数据管道断了,用缓存数据;计算引擎挂了,切到备用实例;分发系统出问题,信号先落地,等恢复后再补发。
总结一下:系统架构没有银弹。我的经验是,先跑通最小闭环,再逐步优化。别一开始就想搞个完美的系统,那往往是最慢的。
嗯,架构这块就聊这么多。记住,再好的策略,没有靠谱的架构支撑,也是空中楼阁。