第二十章:实时监控系统架构:从数据采集到可视化的完整设计
说实话,做量化交易这几年,我踩过最大的坑,不是策略回测过拟合,也不是滑点超预期——而是监控系统在关键时刻挂了。
记得有一次,我正在跑一个高频做市策略,突然市场流动性骤降。我的策略还在傻傻地挂单,结果被一波吃掉。等我发现不对劲,已经亏了六位数。从那以后,我花了大半年时间,重新设计了整套实时监控系统。
今天我就把这套架构拆开给你看。从数据怎么进来,到怎么算,再到怎么展示,一条线讲清楚。
一、整体架构:三层分离,各司其职
我习惯把监控系统分成三层:
- 数据采集层:负责从交易所拿数据
- 计算引擎层:负责算指标、做判断
- 可视化层:负责把结果展示出来
为什么要分层?说白了,就是怕一个模块崩了,整个系统跟着完蛋。你想想看,如果可视化页面卡住了,数据采集还在跑,那至少数据没丢。反过来,如果采集挂了,计算和展示还能报个警。
核心原则:每一层都要能独立运行,层与层之间通过消息队列解耦。
下面这张图是我实际项目中用的架构,你可以参考一下:
二、数据采集:WebSocket 是主力,API 做补充
数据采集这块,我个人的经验是:能用 WebSocket 就别用 REST API。
为什么?因为 WebSocket 是长连接,数据是推过来的,延迟低。REST API 你得轮询,一秒请求几次还好,频率一高,交易所直接给你限流。
我常用的采集方案是这样的:
| 数据类型 | 采集方式 | 频率 | 备注 |
|---|---|---|---|
| 深度行情(Order Book) | WebSocket | 实时推送 | 增量更新,本地维护全量 |
| 成交数据(Trade) | WebSocket | 实时推送 | 用于计算成交量、成交额 |
| K线数据 | REST API | 每1-5秒 | 用于回填和校验 |
| 账户/持仓 | REST API | 每1-10秒 | 频率不宜过高 |
小技巧:WebSocket 连接一定要做心跳检测。我曾经遇到过连接断了但没报错,数据停了半小时才发现。现在我的代码里每5秒检查一次,断了自动重连。
下面是一个简单的 WebSocket 采集代码示例,Python 的:
import websocket
import json
import threading
class MarketDataCollector:
def __init__(self, url, symbols):
self.url = url
self.symbols = symbols
self.ws = None
self.running = False
def on_message(self, ws, message):
data = json.loads(message)
# 这里把数据推到消息队列
self.push_to_queue(data)
def on_error(self, ws, error):
print(f"WebSocket 错误: {error}")
# 自动重连逻辑
self.reconnect()
def on_close(self, ws, close_status_code, close_msg):
print("连接关闭,尝试重连...")
self.reconnect()
def on_open(self, ws):
# 订阅行情
subscribe_msg = {
"type": "subscribe",
"symbols": self.symbols
}
ws.send(json.dumps(subscribe_msg))
def start(self):
self.running = True
self.ws = websocket.WebSocketApp(
self.url,
on_open=self.on_open,
on_message=self.on_message,
on_error=self.on_error,
on_close=self.on_close
)
# 在独立线程运行
wst = threading.Thread(target=self.ws.run_forever)
wst.daemon = True
wst.start()
def reconnect(self):
if self.running:
self.start()
三、计算引擎:Python 搭台,C++ 唱戏
计算引擎这块,我走了不少弯路。一开始全用 Python,后来发现有些计算太慢了。比如流动性深度指标,需要遍历整个订单簿,Python 跑一次要几十毫秒,高频场景根本扛不住。
现在的方案是:Python 做调度和复杂逻辑,C++ 做高频计算。
具体分工是这样的:
- Python 负责:指标计算(流动性比率、买卖压力、波动率等)、异常检测、预警逻辑、数据持久化
- C++ 负责:订单簿快照重建、逐笔成交聚合、延迟敏感的计算(<1ms)
Python 和 C++ 之间怎么通信?我推荐用 ZeroMQ 或者 共享内存。共享内存延迟最低,但实现起来麻烦一点。ZeroMQ 够用,延迟在微秒级。
注意:千万别用 HTTP 做进程间通信。延迟太高,而且容易阻塞。我见过有人用 Flask 做内部接口,结果行情一快,请求排队,整个系统就卡死了。
下面是一个 C++ 计算流动性深度的核心代码片段:
#include <vector>
#include <map>
#include <string>
struct OrderBookLevel {
double price;
double volume;
};
class LiquidityCalculator {
public:
// 计算指定价格区间的流动性深度
double calcDepth(const std::vector<OrderBookLevel>& bids,
const std::vector<OrderBookLevel>& asks,
double mid_price, double depth_pct = 0.01) {
double range = mid_price * depth_pct;
double lower = mid_price - range;
double upper = mid_price + range;
double bid_liquidity = 0.0;
double ask_liquidity = 0.0;
// 累加买单流动性
for (const auto& level : bids) {
if (level.price >= lower) {
bid_liquidity += level.price * level.volume;
}
}
// 累加卖单流动性
for (const auto& level : asks) {
if (level.price <= upper) {
ask_liquidity += level.price * level.volume;
}
}
return bid_liquidity + ask_liquidity;
}
};
四、可视化:Dash + Plotly,够用且灵活
可视化这块,我试过很多方案。Grafana 虽然好看,但定制性差。Bokeh 也不错,但社区不如 Plotly 活跃。最后我选了 Dash + Plotly 的组合。
Dash 的好处是:
- 纯 Python,不用写前端
- 支持实时更新(通过 WebSocket 或轮询)
- 组件丰富,表格、图表、下拉框都有
我一般会展示这几个核心面板:
- 行情概览:最新价、涨跌幅、成交量、买卖盘口
- 流动性仪表盘:深度曲线、买卖压力比、价差变化
- 异常预警:红色闪烁的预警列表,带时间戳和严重等级
- 系统健康:数据延迟、连接状态、CPU/内存使用率
下面是一个 Dash 实时更新的核心代码:
import dash
from dash import dcc, html
from dash.dependencies import Input, Output
import plotly.graph_objs as go
import redis
import json
app = dash.Dash(__name__)
# 从 Redis 读取最新数据
r = redis.Redis(host='localhost', port=6379, db=0)
app.layout = html.Div([
html.H1("流动性监控面板"),
dcc.Graph(id='depth-chart'),
dcc.Interval(id='interval', interval=1000) # 每秒更新
])
@app.callback(
Output('depth-chart', 'figure'),
[Input('interval', 'n_intervals')]
)
def update_depth_chart(n):
# 从 Redis 获取最新深度数据
depth_data = r.get('depth:btc_usdt')
if depth_data is None:
return go.Figure()
data = json.loads(depth_data)
fig = go.Figure()
fig.add_trace(go.Scatter(
x=data['bid_prices'],
y=data['bid_volumes'],
name='买单',
fill='tozeroy'
))
fig.add_trace(go.Scatter(
x=data['ask_prices'],
y=data['ask_volumes'],
name='卖单',
fill='tozeroy'
))
fig.update_layout(title="订单簿深度", xaxis_title="价格", yaxis_title="数量")
return fig
if __name__ == '__main__':
app.run_server(debug=True, host='0.0.0.0', port=8050)
五、避坑指南:我踩过的几个坑
最后,分享几个我实际项目中遇到的坑,希望能帮你少走弯路:
坑1:数据时间戳不一致
不同交易所的时间戳格式不一样,有的用毫秒,有的用微秒。我一开始没注意,结果两个交易所的数据对不上,算出来的指标全是错的。解决方案:统一转成 UTC 毫秒时间戳,在采集层就做归一化。
坑2:消息队列积压
行情一快,Kafka 消费不过来,消息越积越多。我后来加了背压机制:如果队列长度超过阈值,就丢弃非关键数据(比如中间价变化),只保留关键数据(比如深度快照)。
坑3:可视化页面卡死
Dash 默认是同步更新,如果数据量大,页面会卡。我后来改成异步回调,并且限制图表数据点数量(最多保留500个点),超过就丢弃旧数据。
嗯,这套架构我用了快两年,稳定性还不错。当然,每个交易场景不一样,你可以根据自己的需求调整。比如如果你只做低频,Python 全栈就够了,不用上 C++。但如果你做高频,那 C++ 这层是省不了的。
记住一句话:监控系统不是锦上添花,是保命用的。花时间把这块做好,绝对值。