Who notices when a price feed goes wrong:行情源异常监控告警系统落地实践
做量化交易、行情展示或者金融数据服务的时候,最怕的不是行情跌了,而是行情数据源悄悄出了问题,系统还在“正常”运行。这个问题听起来很简单:“谁会注意到价格源错了?”但在真实工程里,答案往往让人头疼——大多数情况下,第一个发现异常的不是系统,而是交易员或者下游业务方,等到人眼看出价格不对劲,可能已经过去了十几分钟甚至更久。
本文围绕“如何让系统第一时间发现行情源异常”展开,从数据源故障类型、监控分层设计、核心检测算法,到 Python 实现、告警通知、值班响应和工程化最佳实践,完整拆解一套可用于实际项目的行情源监控告警方案。无论你是量化开发、后端工程师,还是负责行情源接入的数据开发,都可以照着落地。
1. 背景:为什么“价格源出错”很难被及时发现
1.1 什么是 Price Feed
Price Feed 通常翻译为“行情源”或“价格源”,指的是一类持续输出标的价格数据的系统或接口。在量化交易里,它可能是交易所的撮合行情 WebSocket;在金融信息服务里,它可能是授权商提供的实时快照;在区块链链上服务里,它可能是预言机节点推送的价格。
行情源的核心使命是:持续、稳定、准确地提供某个品种的价格。它的下游可能是策略引擎、止损模块、账本估值系统,甚至只是一个手机 App 上的 K 线图。
1.2 行情源出错往往不是“断网”这么简单
如果行情源直接断开,系统报错,这个问题反而容易发现。真正危险的是“半坏”状态:
- 数据没有完全断,但更新频率从 1 秒一次变成了 30 秒一次;
- 某个交易对的价格突然跳到一个明显不合理的位置;
- 上游做市商维护,价格数据停滞在收盘价;
- 上游返回的时间戳仍旧在跳动,但背后的价格已经被冻结;
- 不同数据源之间的价格出现明显背离,而且持续了很长时间。
这些问题在数据源内部不一定被判定为故障,但在下游业务眼里,它们已经足以造成策略误判。一个经典的例子是:止损模块依赖行情源判断是否触发风控,如果价格源滞后几分钟,一旦市场剧烈波动,止损单可能根本来不及触发。
1.3 “Who notices”在工程中的真实含义
从工程视角看,标题里的 “Who notices” 不是在问某一个具体的人,而是在问:系统的可观测性设计,是否能让“最先发现问题”的角色从“人肉巡检”变成“自动监控”?
我们想要达到的状态是:
- 机器第一时间计算行情源的健康度;
- 监控系统自动定位是哪类异常(延迟、跳变、偏差);
- 告警通知到正确的人;
- 值班人员根据预案快速处理或切换数据源。
这件事不能靠偶然,需要一套完整的代码、配置和响应流程。
2. 行情源监控的核心指标与异常类型
在设计监控系统之前,先需要定义“出错”怎么量化。根据常见的生产故障,可以把行情源异常归纳为以下几类。
2.1 新鲜度异常(Stale / No Update)
行情源在预期时间内没有产生新数据。
量化指标:
last_update_time:最近一次收到行情数据的时间;expected_interval:正常情况下行情数据的最大间隔。
判定逻辑:
当前时间 - 最近一次收到数据的时间 > 阈值这里的阈值要根据品种特性设定。股票逐笔可能要求 1 秒级,合约深度行情可能 500ms 算正常,而某些日线快照 24 小时更新一次就足够。阈值不能写死,应该做成配置。
2.2 价格跳变异常(Price Spike)
相邻两个价格之间的变化幅度超出合理范围。
量化指标:
previous_price:上一次价格;current_price:当前价格;max_change_ratio:允许的最大变化比例。
判定逻辑:
abs(current_price - previous_price) / previous_price > max_change_ratio但这里有个误区:真实市场里确实会出现瞬间大幅波动,尤其是流动性差的品种或者新闻驱动行情。因此跳变检测一般只作为“预警”,需要结合偏差检测或者人工确认,不能只看一个指标下结论。
2.3 多源偏差异常(Deviation)
如果系统同时接入了多个行情源(例如交易所A、交易所B、聚合行情源C),可以通过交叉对比发现某个源是否异常。
量化指标:
price_source_a:A 源的当前价格;price_source_b:B 源的当前价格;max_deviation_ratio:允许的最大偏差比例。
这个方案在实际项目里非常有效。正常情况下,即使不同交易所之间存在价差,也会维持在一个较小区间内。一旦某个源异常,它的价格会短时间内与其他源严重脱节。
2.4 时间戳异常(Timestamp Out of Sync)
部分行情源在数据包里带有自己的时间戳,这时候需要检查上游时间与本地服务器时间是否在合理范围内。
常见场景是:
- 上游服务器时钟出现偏差;
- 推送消息内部的时间戳字段没有正确更新;
- 本地消费端引入消息队列后出现时间解析错误。
2.5 残缺数据与精度异常
字符串格式数据源容易在解析阶段失败。这类异常通常表现为:
- 返回消息缺少必填字段;
- 价格被错误地除以或乘以 100/1000;
- 小数位发生变化,例如从 1.234 变成 123.4。
这类错误最难发现,因为在数据格式上它是合法的,只有在与其他源对比时才会暴露。
3. 监控系统总体架构分层
一个完整的行情源监控告警系统,可以拆成四个层次。
3.1 采集层(Collector)
负责从行情源读取数据,统一转换为内部标准结构。采集层不关心数据对不对,只负责“收数据”和“打点”。
核心逻辑:
- 连接行情源并订阅品种;
- 解析原始报文并转换为统一的
MarketData结构; - 每条数据写入内存缓存、消息队列或时序数据库;
- 上报心跳时间到监控模块。
3.2 校验层(Validator)
校验层是监控的核心。每收到一条行情数据,都执行一组检测规则,例如:
- 延迟是否加大;
- 价格是否跳变;
- 与其他源相比是否偏离;
- 更新时间是否推进。
校验层输出“健康检查结果”,而不是直接推送告警,这样便于上层灵活决定通知策略。
3.3 告警层(Alert Manager)
告警层决定要不要通知、通知谁、怎么通知。它负责:
- 对同一品种的重复告警进行聚合;
- 根据级别匹配不同通知渠道;
- 记录告警的触发时间和恢复时间;
- 在长时间未恢复时将告警升级。
3.4 响应层(Responder)
响应层由值班流程组成:
- 收到告警后,先判断是否为真实故障;
- 如果是上游故障,切换到备用数据源;
- 如果短时间无法恢复,通知交易或风控人员冻结相关策略;
- 故障处理后,确认数据恢复正常恢复交易。
整体流程可以概括为:
行情数据源 ↓ 采集层订阅推送 ↓ 写入缓存 + 记录心跳 ↓ 校验层执行健康规则 ↓ 触发异常时生成告警事件 ↓ 告警层聚合去重并推送给值班人员 ↓ 值班人员响应处理 / 自动切换备用源下面我们用 Python 实现一个相对完整的监控核心逻辑。
4. 核心代码实现:让你的系统“第一个发现问题”
下面代码以 Python 为例,目标是实现一个行情源监控守护进程。它不绑定具体行情商,而是抽象出两个关键接口:
- 行情数据来源(模拟真实推送);
- 告警通知器(模拟发送 webhook)。
你可以在此基础上替换为真实的数据源 SDK 和通知机器人。
4.1 项目结构
price_feed_monitor/ ├── config.py # 全局配置 ├── models.py # 数据模型 ├── collectors.py # 行情源采集适配层 ├── validators.py # 异常检测规则 ├── notifiers.py # 告警通知 ├── monitor.py # 主监控循环 └── requirements.txt # 依赖# requirements.txt requests>=2.31.0本文示例不依赖定时框架,采用time.sleep轮询模式,便于理解核心逻辑。实际生产环境可以换成 APScheduler 或基于异步事件循环。
4.2 全局配置 config.py
# config.py # 各品种的监控阈值(单位:秒 或 比例) RULES = { "BTC-USDT": { "max_stale_seconds": 10, # 超过10秒没有新行情就告警 "max_price_jump_ratio": 0.01, # 相邻价格跳变超过1%告警 "max_deviation_ratio": 0.005, # 与参考源偏差超过0.5%告警 }, "ETH-USDT": { "max_stale_seconds": 10, "max_price_jump_ratio": 0.015, "max_deviation_ratio": 0.008, }, } # 全局检查周期,单位秒 CHECK_INTERVAL = 1 # 备选数据源 URL(示例中只用于模拟) DATA_SOURCE_CONFIG = { "primary": "ws://your-price-source.example.com/market", "backup": "ws://backup.example.com/market", }不同品种的波动特性不同,因此阈值必须做成按品种配置,而不是全局一个固定值。新建交易对时,先观察一段时间的历史波动,再确定合理阈值。
4.3 数据模型 models.py
# models.py from dataclasses import dataclass from datetime import datetime @dataclass class MarketData: """统一行情数据模型""" symbol: str # 交易对,例如 BTC-USDT price: float # 最新价 volume: float # 成交量(可选) source: str # 数据源标识 received_at: datetime # 本地接收到数据的时间 upstream_time: datetime = None # 上游自带时间戳,可选 @dataclass class HealthStatus: """某个品种的健康状态""" symbol: str is_healthy: bool error_code: str = "" # STALE / PRICE_SPIKE / DEVIATION message: str = "" updated_at: datetime = None使用dataclass是为了让代码清晰,并且方便后续序列化成 JSON。
4.4 采集适配层 collectors.py
真实项目中,采集层会连接 WebSocket 或 API。在下方案例中,我们用一个“模拟行情源”来演示,它按固定间隔产生新的行情。当你接入真实环境时,替换generate_market_data函数即可。
# collectors.py import random import time from datetime import datetime from threading import Thread, Event from models import MarketData class BaseCollector: """行情源采集器基类""" def __init__(self, source_name: str, symbols: list[str]): self.source_name = source_name self.symbols = symbols self._stop_event = Event() def start(self): raise NotImplementedError def stop(self): self._stop_event.set() class SimulatedCollector(BaseCollector): """ 模拟行情源。 正常情况下每 1 秒推送一次价格。 为了演示故障,可以设置 stale_after_seconds 模拟数据停更。 """ def __init__( self, source_name: str, symbols: list[str], on_data, stale_after_seconds: int = None, base_price: float = 60000.0, ): super().__init__(source_name, symbols) self._on_data = on_data self._stale_after_seconds = stale_after_seconds self._base_price = base_price def start(self): print(f"[Collector] {self.source_name} 开始运行") thread = Thread(target=self._run_loop, daemon=True) thread.start() def _run_loop(self): count = 0 while not self._stop_event.is_set(): # 模拟数据停更故障 if self._stale_after_seconds and count >= self._stale_after_seconds: time.sleep(1) continue for symbol in self.symbols: # 随机小幅波动 price = self._base_price + random.uniform(-5, 5) data = MarketData( symbol=symbol, price=round(price, 2), source=self.source_name, received_at=datetime.now(), ) self._on_data(data) count += 1 time.sleep(1)这段代码的关键作用不是模拟数据本身,而是提供一个可以替换的数据接入点。on_data回调就是后面校验模块的入口。
4.5 异常检测规则 validators.py
# validators.py import time from collections import defaultdict from datetime import datetime from models import HealthStatus, MarketData class PriceFeedMonitor: """ 行情监控核心:维护最新快照,执行多种异常检测。 """ def __init__(self, rules: dict, notifier): self.rules = rules self.notifier = notifier # 保存每个源、每个品种的最新行情快照 self._latest_data = {} # 保存最近一次告警时间,用于去重 self._last_alert_time = defaultdict(float) def on_market_data(self, data: MarketData): """每条新行情进来都会调用这里""" key = f"{data.source}:{data.symbol}" self._latest_data[key] = data statuses = [ self._check_stale(data), self._check_price_jump(data), self._check_deviation(data), ] for status in statuses: if status and not status.is_healthy: self._trigger_alert(status) def _check_stale(self, data: MarketData): """ 新鲜度检测: 如果当前时间距最近一条数据时间过长,说明行情源可能停更。 """ if data.symbol not in self.rules: return None rule = self.rules[data.symbol] max_stale = rule.get("max_stale_seconds", 10) time_diff = (datetime.now() - data.received_at).total_seconds() if time_diff > max_stale: return HealthStatus( symbol=data.symbol, is_healthy=False, error_code="STALE", message=f"{data.source} 超过 {max_stale} 秒未更新行情", updated_at=datetime.now(), ) return None def _check_price_jump(self, data: MarketData): """ 跳变检测:比较同一数据源自身的前后价格差异。 这里需要对比相同 source 的上一笔价格。 """ rule = self.rules.get(data.symbol) if not rule: return None key = f"{data.source}:{data.symbol}" history_key = f"history:{key}" last_price = self._latest_data.get(key) if not last_price or last_price.price == 0: return None ratio = abs(data.price - last_price.price) / last_price.price max_jump = rule.get("max_price_jump_ratio", 0.01) if ratio > max_jump: return HealthStatus( symbol=data.symbol, is_healthy=False, error_code="PRICE_SPIKE", message=( f"{data.source} 价格跳变 {ratio:.4%}," f"前值 {last_price.price},当前 {data.price}" ), updated_at=datetime.now(), ) return None def _check_deviation(self, data: MarketData): """ 偏差检测:需要在多个数据源之间进行对比。 当前示例只得到同一源的最新数据,这里演示单源模式。 真实生产可维护 source1/source2 两个快照进行 cross-check。 """ return None def _trigger_alert(self, status: HealthStatus): """告警去重:同一品种、同一错误码 10 秒内不重复推送""" alert_key = f"{status.symbol}:{status.error_code}" now = time.time() if now - self._last_alert_time.get(alert_key, 0) < 10: return self._last_alert_time[alert_key] = now self.notifier.send(status)这里有几个可以优化的点:
- 跳变检测需要考虑“多次连续触发”的情况,避免由单个报点异常导致误报。
- 新鲜度检测最好由独立的调度任务周期执行,而不是只在收到行情时执行。因为一旦行情源完全停更,
on_market_data永远不会再被调用。 - 多源偏差检测事实上需要独立维护多个 source 的快照,建议用一个字典维护
symbol -> {source1: price, source2: price}。
为了弥补上面提到的“停更后不再有回调”的问题,我们需要在monitor.py中启动一个独立线程周期检查新鲜度。
4.6 告警通知 notifiers.py
通知模块支持两个实现:
- 控制台输出;
- 通用 Webhook 机器人(企业微信、钉钉、飞书都常见)。
# notifiers.py import json import requests from models import HealthStatus class ConsoleNotifier: def send(self, status: HealthStatus): print(f"[ALERT] {status.error_code} | {status.symbol} | {status.message}") class WebhookNotifier: """ 使用通用 Webhook 发送告警,可根据实际平台调整 JSON 模板。 """ def __init__(self, webhook_url: str): self.webhook_url = webhook_url def send(self, status: HealthStatus): payload = { "msgtype": "text", "text": { "content": ( f"行情源异常告警\n" f"品种:{status.symbol}\n" f"类型:{status.error_code}\n" f"详情:{status.message}" ) } } try: resp = requests.post( self.webhook_url, json=payload, timeout=5 ) resp.raise_for_status() print(f"[Notifier] 已推送告警到 webhook,状态码 {resp.status_code}") except Exception as exc: print(f"[Notifier] 推送告警失败:{exc}")Webhook 的 JSON 模板各平台不同,实际使用时请以机器人平台文档为准,这里只展示通用思路。
4.7 独立守护监控任务 monitor.py
# monitor.py import threading import time from datetime import datetime from collectors import SimulatedCollector from models import HealthStatus from notifiers import ConsoleNotifier, WebhookNotifier from validators import PriceFeedMonitor def create_demo_monitor(): # 这里使用控制台通知,便于观察;真实环境请改成 WebhookNotifier notifier = ConsoleNotifier() monitor = PriceFeedMonitor(rules=RULES, notifier=notifier) collector = SimulatedCollector( source_name="demo_source", symbols=["BTC-USDT", "ETH-USDT"], on_data=monitor.on_market_data, # 模拟 8 秒后数据停更 stale_after_seconds=8, base_price=60000.0, ) return monitor, collector def periodic_stale_check(monitor: PriceFeedMonitor): """ 周期检查所有已注册行情源的新鲜度。 因为行情源完全停止推送时,on_market_data 不会再触发, 所以必须有一个后台任务主动巡查。 """ while True: now = datetime.now() latest_items = list(monitor._latest_data.items()) for key, data in latest_items: rule = RULES.get(data.symbol) if not rule: continue max_stale = rule.get("max_stale_seconds", 10) time_diff = (now - data.received_at).total_seconds() if time_diff > max_stale: status = HealthStatus( symbol=data.symbol, is_healthy=False, error_code="STALE", message=f"{data.source} 后台巡检发现超过 {max_stale} 秒未更新", updated_at=now, ) monitor._trigger_alert(status) time.sleep(1) if __name__ == "__main__": from config import RULES monitor, collector = create_demo_monitor() collector.start() # 启动后台巡检线程 thread = threading.Thread(target=periodic_stale_check, args=(monitor,), daemon=True) thread.start() print("行情源监控已启动,按 Ctrl+C 停止。") try: while True: time.sleep(1) except KeyboardInterrupt: collector.stop() print("监控已停止")4.8 运行与预期输出
直接运行:
python monitor.py如果设置stale_after_seconds=8,可以观察到类似输出:
[Collector] demo_source 开始运行 [ALERT] STALE | BTC-USDT | demo_source 后台巡检发现超过 10 秒未更新 [ALERT] STALE | ETH-USDT | demo_source 后台巡检发现超过 10 秒未更新由于告警做了 10 秒去重,后续不会因为同一原因反复刷屏。等到数据源恢复推送后,系统会继续正常接收新的行情,此时可以再补一个“恢复通知”。
5. 从单机脚本到工程级方案的进阶要点
上面的脚本能解决“发现异常”的问题,但离生产环境还有差距。下面列出从 Demo 升级到工程级方案时需要关注的模块。
5.1 恢复通知
很多团队只实现了“故障告警”,却忘了恢复通知。结果就是值班人员处理完故障后,还需要手动确认数据是否已经恢复。建议在状态从异常变为正常时,推送一条恢复通知。
实现方式:
def _check_and_notify_recovery(self, status: HealthStatus): # 记录状态机,上次异常、本次正常时触发恢复通知 pass5.2 多源交叉比对
如果系统接入了两个以上行情源,最好能实现真正的多源交叉比对。实现在validators.py中新增:
def _check_cross_source_deviation(self, symbol: str): prices = {} for key, data in self._latest_data.items(): if data.symbol == symbol: prices[key.split(":")[0]] = data.price if len(prices) < 2: return None # 以第一个源为基准,计算其他源与基准的偏差 source_names = list(prices.keys()) base_source = source_names[0] base_price = prices[base_source] for src in source_names[1:]: deviation = abs(prices[src] - base_price) / base_price if deviation > self.rules[symbol]["max_deviation_ratio"]: # 生成偏差告警,标记为 DEVIATION pass注意,基准源本身也可能异常,更稳妥的做法是取多个源的中位数作为参考价格,再计算每个源与参考价格的偏差。
5.3 存储与历史回溯
监控产生的事件最好写入时序数据库,例如 InfluxDB、Prometheus 或 ClickHouse,方便后续:
- 复盘故障发生时间;
- 统计数据源可用率 SLA;
- 回溯异常告警与策略交易时间的关系。
如果团队暂未引入时序数据库,也可以先写入 MySQL、PostgreSQL 或 Elasticsearch。关键在于:告警不只是即时消息,也是留给未来的审计资料。
5.4 告警升级机制
长时间未恢复的异常需要升级到更高级别的负责人。常见策略:
- 告警触发后 5 分钟未恢复,通知数据组;
- 15 分钟未恢复,通知交易组和风控组;
- 30 分钟未恢复,自动禁用故障行情源并切换备用源。
这里要结合团队实际,不能盲目自动化。特别是在策略交易场景,禁用行情源可能影响在途订单,需要人工确认确认后再执行。
6. 常见问题与排查思路
| 问题现象 | 常见原因 | 解决思路 |
|---|---|---|
| 行情源已经停更,但监控没有告警 | 新鲜度检测依赖on_market_data,行情停更后不会触发回调 | 增加独立后台巡检线程,周期扫描所有快照的时间戳 |
| 行情偶尔出现一次大跳变就疯狂告警 | 跳变阈值设置过小,或单个异常报点没有经过多轮确认 | 提高阈值;采用“连续 N 次触发才告警”策略 |
| 不同数据源正常价差就大于告警阈值 | 阈值设置不合理,没有结合真实市场数据 | 收集历史价差数据,按 99 分位数设定阈值 |
| 告警重复轰炸,值班人员免疫 | 缺少去重机制 | 同一品种同一错误码在时间窗口内只推送一次 |
| 故障恢复后仍然在告警 | 没有实现恢复通知和状态机 | 维护每品种每数据源的状态,异常切正常时推送恢复 |
| 多个监控脚本重复部署,告警重复 | 没有统一的任务注册中心 | 使用独立的任务调度服务,或在存储层做唯一键去重 |
| Webhook 推送失败没有感知 | 通知模块异常被吞掉 | 记录通知失败日志,并对通知渠道本身做兜底监控 |
7. 最佳实践与工程化建议
7.1 先定义 SLA,再写监控规则
没有 SLA,就不知道阈值该怎么定。可以先定义:
- 行情数据最大延迟不超过 N 秒;
- 数据源月度可用率不低于 99.9%;
- 告警响应时间不超过 5 分钟。
这些目标决定了阈值、巡检频率和值班流程。建议先跑 1 到 2 周数据,再调整规则,刚开始阈值不要定太紧。
7.2 不要只监控价格,还要监控数据全链路
行情从上游到策略进程,通常经过好几个环节:
上游行情源 → 采集服务 → 消息队列 → 数据加工 → 策略进程任何一个环节卡住,都会导致下游看不到新行情。建议在每一条链路都输出心跳指标,例如:
- 采集服务收到消息后打点
collector_received_count; - 消费端每消费一条后打点
consumer_processed_count; - 策略进程记录最新行情时间戳。
只有监控到“最终策略最后收到的行情时间”,才能确认端到端健康。
7.3 多源交叉校验是性价比最高的方案
单一数据源的“自检”很难发现逻辑错误,例如精度变了、价格字段错位、交易所维护时价格冻结。交叉校验通过不同数据源之间的对比,能发现绝大多数这类问题。
但要注意,两个数据源都错的情况也不是没有,例如某些小币种流动性不足,行情源价格本身就失真。遇到这种情况,需要人工复盘或引入第三方参考源。
7.4 故障演练要常态化
不要等到真实故障时才验证监控是否有效。建议定期演练:
- 手动停掉主行情源的推送;
- 观察监控报警时间是否符合预期;
- 验证备用源切换流程能否跑通;
- 记录整体恢复时间,评估是否符合 SLA。
演练可以按季度执行,同时保留完整演练报告。
7.5 数据合规与授权提醒
接入任何行情数据源之前,先确认授权范围。不同市场、不同数据商对行情数据的转发、存储、展示有严格限制。量化交易项目尤其需要注意:即使只是内部估值使用,也不能随意抓取未经授权的行情源。工程同学不要为了省事引入来源不明的数据,合规风险远大于技术收益。
7.6 监控系统自身的“高可用”
如果监控服务只部署一个实例,它挂了怎么办?
基础做法:
- 监控服务以守护进程方式启动,由 systemd 或常驻进程管理器托管;
- 告警接口使用外部平台(企业微信、钉钉、飞书、邮件),避免和监控服务共用一套进程;
- 保存最近 N 条行情数据在本地文件或 Redis,方便排查时快速查看。
进阶做法:
- 部署双机互备或主从模式;
- 使用独立的监控探针,检查主监控进程是否存活。
7.7 状态机与降噪
在复杂业务中,告警要经历“原始告警 → 预确认 → 正式告警 → 恢复”的完整生命周期。可以把状态集中管理:
| 错误码 | 触发条件 | 预确认策略 | 处理动作 |
|---|---|---|---|
| STALE | 超过最大间隔未更新 | 连续 3 次巡检均未恢复 | 通知数据负责人 |
| PRICE_SPIKE | 单次价格跳变超过阈值 | 连续 2 次推送跳变或对比备用源 | 通知交易风控 |
| DEVIATION | 多源偏差超过阈值 | 基准源取多源中位数 | 通知数据组,暂停使用异常源 |
| TIMESTAMP | 时间戳与本地差异过大 | 检测到一次即告警 | 检查上游时间解析逻辑 |
通过这种方式,把“收到异常”和“确认故障”分开,能显著减少值班人员的噪声负担。
8. 回到标题:究竟谁来“notice”
行情源是一个低频调用的基础设施,越稳定,越容易被忽略。等到你真正发现它出错时,往往已经出现了可感知的损失。
本文给出的思路并不复杂:把“靠人发现问题”改造成“靠系统发现问题”。具体来说,就是做好四件事:
- 定义哪些现象算异常,给每个品种配置合理阈值;
- 在采集链路中自动执行新鲜度、跳变、多源偏差检查;
- 将异常转化为分级告警,并通过 Webhook 通知到人;
- 通过状态机、去重、恢复通知和巡检任务,保证监控自身不成为新的盲点。
这一整套能力并不需要庞大的系统,一个 Python 脚本也可以作为起点。关键在于:当行情源坏了的第一秒,系统就应该知道,并且让该知道的人知道。
如果这篇文章对你有帮助,可以先按文中代码在本地跑通一个最小 Demo,再逐步加入你自己的行情源、消息队列和值班流程。如果只是看懂了却不动手,下次行情源出错时,第一个发现的可能真的还是交易员。