第二十章:实时监控系统架构:从数据采集到可视化的完整设计

说实话,做量化交易这几年,我踩过最大的坑,不是策略回测过拟合,也不是滑点超预期——而是监控系统在关键时刻挂了。

记得有一次,我正在跑一个高频做市策略,突然市场流动性骤降。我的策略还在傻傻地挂单,结果被一波吃掉。等我发现不对劲,已经亏了六位数。从那以后,我花了大半年时间,重新设计了整套实时监控系统。

今天我就把这套架构拆开给你看。从数据怎么进来,到怎么算,再到怎么展示,一条线讲清楚。

一、整体架构:三层分离,各司其职

我习惯把监控系统分成三层:

  • 数据采集层:负责从交易所拿数据
  • 计算引擎层:负责算指标、做判断
  • 可视化层:负责把结果展示出来

为什么要分层?说白了,就是怕一个模块崩了,整个系统跟着完蛋。你想想看,如果可视化页面卡住了,数据采集还在跑,那至少数据没丢。反过来,如果采集挂了,计算和展示还能报个警。

核心原则:每一层都要能独立运行,层与层之间通过消息队列解耦。

下面这张图是我实际项目中用的架构,你可以参考一下:

实时监控系统三层架构 数据采集层 WebSocket 行情 REST API 快照 订单流数据 其他源 消息队列(Kafka / Redis) 计算引擎层 Python 指标计算 C++ 高频处理 异常检测引擎 预警模块 可视化层(Dash / Plotly)

二、数据采集: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 或轮询)
  • 组件丰富,表格、图表、下拉框都有

我一般会展示这几个核心面板:

  1. 行情概览:最新价、涨跌幅、成交量、买卖盘口
  2. 流动性仪表盘:深度曲线、买卖压力比、价差变化
  3. 异常预警:红色闪烁的预警列表,带时间戳和严重等级
  4. 系统健康:数据延迟、连接状态、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++ 这层是省不了的。

记住一句话:监控系统不是锦上添花,是保命用的。花时间把这块做好,绝对值。