从一次实际需求说起。年前我在给一个量化项目做数据链路改造手里有一个 QMT 的 Python 客户端正在接收实时行情。刚开始一切正常策略进程自己订阅自己消费。但很快出现了一个新的需求另一个 Web 服务也要展示实时行情还有一台监控面板也想看到行情。于是问题不再是“怎么接行情”而是“怎么把行情从 QMT 的进程里安全地分发给多个进程”。这就是跨进程行情桥的典型场景。而这次我选择的方式是先用 UDP 把行情广播出去再由一个桥接进程把 UDP 转成 WebSocket供前端和跨进程服务消费。一句话说清主判断这个桥不是“从一种协议换成另一种协议”而是把“单进程里的实时订阅”升级成“多进程可共享的行情通道”。协议只是表面真正改变的是数据的所有权和分发边界。1. 先搞清楚桥要解决的是哪一个真问题很多人一看到“UDP 转 WebSocket”就以为这是个协议转换项目。协议转换当然是一个动作但如果你只把它当协议转换来做做完会发现自己既没有解决性能问题也没有解决架构问题只是把报文格式从二进制换成了 JSON 而已。1.1 一个常见的 QMT 行情分发困境QMT 在量化交易里的主要工作方式是通过 Python API 在本进程内完成行情订阅和回调。数据在回调函数里谁写代码谁用没问题。但当你出现下面的场景时麻烦就来了策略进程要消费行情同时有一个 Web 前端要展示行情。你有多个策略进程各自都要订阅行情但 QMT 的账号和连接可能只适合一个主进程维护。你不想把策略的核心逻辑和行情转发逻辑耦合在一起希望行情先“出来”再被不同模块取用。你希望行情不只留在本机而是可以给局域网内其他服务消费。这时候你会意识到QMT 的行情接口天然是“进程内”的。你需要在进程边界上开一个口子让行情数据可以流出去。这就是跨进程行情桥的起点。1.2 协议选择的底层逻辑UDP 的语义与 WebSocket 的语义为什么中间要过一道 UDP而不是直接在 QMT 进程里开 WebSocket 服务第一个原因是职责分离。让 QMT 客户端进程既管行情订阅又管 WebSocket 服务端不是不行但会带来耦合WebSocket 客户端的接入、断开、心跳、广播逻辑会把行情处理流程越拉越长。先把行情通过 UDP 打出去再由独立进程做 WebSocket 服务两边可以独立演进、独立重启。这件事本质上类似“进程间的消息总线”。第二个原因是 QMT 行情本身的推送语义。行情数据是高频、单向、事务性的推送流和 UDP 的语义很匹配发送端只管发不关心谁在听接收端自己处理迟到、丢失和重复。如果你强行在 QMT 进程里用 TCP 或 WebSocket 直接推送会面临连接管理、慢消费者、背压等一系列问题这些都不该由行情采集进程来承担。第三个原因是 WebSocket 的消费语义。WebSocket 是连接型的谁连上来、推给谁、断开了怎么重连这些都很清晰。它对 Web 前端、服务端中间件、跨语言客户端都非常友好。所以合理的分工是UDP 负责本地或局域网内的低延迟广播WebSocket 负责面向多消费者的可靠连接通道。这里有一个容易被忽略的点UDP 转 WebSocket 并不是一个“谁更好”的选择而是把两种不同语义的协议放到各自最合适的位置。UDP 管发送WebSocket 管分发。2. 桥的总体架构一条三节点链路整个链路可以分为三个节点理解这三个节点的职责比理解代码本身更重要。2.1 节点一QMT 行情客户端UDP 发送端这个节点是一个 Python 进程内部使用 QMT 的行情接口订阅行情。它在回调函数里拿到行情数据后不是直接打印或写入某个数据结构而是立刻序列化并通过 UDP Socket 发送出去。常见的发送方式有两种单播发送到127.0.0.1:12000或局域网内固定 IP适合只有单一桥接进程的场景。广播发送到255.255.255.255:12000适合同一网段有多台机器需要接收的场景。我一般建议先用单播。原因是广播会带来不必要的网络压力而且会让排查链路多出一个变量你不知道数据是被哪台机器消费了还是被某台机器的防火墙挡了。发送端要解决的问题包括行情数据的高频拼包。序列化格式的统一。发送频率和 UDP 缓冲区的关系。进程退出时如何释放 Socket。示意的 Python 发送端并不复杂重点是把回调函数里的数据转成字节流import socket import json import time udp_sender socket.socket(socket.AF_INET, socket.SOCK_DGRAM) udp_sender.connect((127.0.0.1, 12000)) def on_bar(data): payload json.dumps(data).encode(utf-8) udp_sender.send(payload)这段代码只是示意结构真实项目里你还需要处理行情字段是否需要裁剪。发送频率是否需要限速。数据是否需要加一个序号方便接收端判断乱序和丢包。2.2 节点二UDP 到 WebSocket 的桥接进程这是核心节点。它监听 UDP 端口收到数据包后解析然后通过 WebSocket 推送给所有已连接的客户端。桥接进程既是一个 UDP Server也是一个 WebSocket Server。它在两个协议之间做消息转换和状态管理。推荐用 Python 来做因为实现成本低而且 websockets 库足够成熟。核心结构是启动 UDP 监听线程持续接收行情包。解析 JSON 后写入一个内部队列。WebSocket Server 在另一个事件循环里运行一旦有客户端连接就把队列中的行情推送出去。多个客户端同时在线时采用广播策略所有客户端共享同一份行情流。这里要注意一个比较隐蔽的设计问题UDP 接收速度和 WebSocket 推送速度不一定匹配。如果大量行情瞬间涌入而某个 WebSocket 客户端处理不过来你是丢弃这部分行情还是为这个客户端做消息积压参考答案是默认丢弃或按客户端订阅做过滤。行情这类实时数据过期数据没有回补价值积压只会让延迟越拉越大。如果你需要可靠送达应该走另一套基于消息队列的方案而不是在这个桥里硬扛。import asyncio import json import socket import websockets UDP_IP 127.0.0.1 UDP_PORT 12000 WS_PORT 13000 connected_clients set() async def udp_listener(): loop asyncio.get_running_loop() udp_sock socket.socket(socket.AF_INET, socket.SOCK_DGRAM) udp_sock.bind((UDP_IP, UDP_PORT)) udp_sock.setblocking(False) while True: data, addr await loop.sock_recvfrom(udp_sock, 65535) message json.loads(data.decode(utf-8)) if connected_clients: await asyncio.gather( *[client.send(json.dumps(message)) for client in connected_clients], return_exceptionsTrue ) async def ws_handler(websocket): connected_clients.add(websocket) try: await websocket.wait_closed() finally: connected_clients.discard(websocket) async def main(): async with websockets.serve(ws_handler, 0.0.0.0, WS_PORT): await udp_listener() asyncio.run(main())这段代码只解决了“能跑”不能直接上生产。真实工程里还需要加上心跳、断线清理、客户端数量限制、消息过滤、日志和异常隔离。但骨架就是这样一个 UDP 收包一组 WebSocket 连接一个广播循环。2.3 节点三WebSocket 下游消费者下游消费者可以是任何支持 WebSocket 的客户端浏览器前端用new WebSocket(ws://localhost:13000)接收行情。Python 事件驱动服务用websockets客户端库持续消费。其他语言的服务比如 Go 或 Node.js通过标准 WebSocket 协议接入。因为桥接进程把 UPD 报文统一转成了 JSON下游消费者不需要关心 QMT 的数据结构细节拿到什么字段直接渲染或入库存即可。这里我建议在消息里加一个字段msg_type用来区分行情类型、订阅确认、心跳和系统错误。否则下游消费者很难区分一条消息是行情还是心跳。{ msg_type: quote, symbol: 600000.SH, price: 10.25, ts: 1730000000 }这样一个简单结构足以支撑前端实时刷新、后端指标计算和日志审计。3. 实现细节从最小跑通到稳定运行只讲架构太虚只讲代码太窄。这一节把实现路径拆成三步每一步都先跑通再优化。3.1 环境准备与消息格式设计建议环境Python 3.9 以上。websockets库安装方式QMT Python API根据你使用的版本安装到对应环境里注意 32 位和 64 位的区别。消息格式是整个桥的契约。建议先定义好再动手写代码。我的建议是把行情消息拆成两层外层是通用信封包含msg_type、ts、source。内层是具体行情字段包含标的代码、最新价、成交量、时间戳等。不要一开始就设计太复杂的嵌套结构因为桥的核心是“低延迟转发”字段越简单越好。等业务需要时再在消费端扩展。3.2 桥接进程的关键代码思路队列、广播与异常隔离桥接进程里最容易出错的地方不在 UDP 收包也不在 WebSocket 发送而在两个并发循环之间的资源竞争。UDP 接收线程是高频的WebSocket 事件循环也是高频的。如果直接把socket.recvfrom放进异步事件循环可能会导致 UDP 在高负载下阻塞整个 WebSocket Server。所以更好的做法是让 UDP 接收保持在独立的线程或独立事件循环里用线程安全队列把消息传给 WebSocket 推送器。实际工程中我建议使用一个有界队列from queue import Queue message_queue Queue(maxsize10000)队列满了就丢弃新消息而不是阻塞 UDP 接收。记录丢弃次数方便后续对行情量和消费能力做评估。如果 UDP 端的消息量实在太快比起硬撑队列更好的办法是做行情裁剪或降低订阅数。WebSocket 发送时也要做异常隔离。某个客户端断开或异常慢不能影响其他客户端。比较稳妥的方式是给每个连接一个单独的发送缓冲推送失败时只关闭该连接async def safe_send(websocket, message): try: await websocket.send(message) except Exception: connected_clients.discard(websocket) await websocket.close()3.3 跨进程验证先本地再局域网很多人在本地就能跑通一上局域网就出问题。原因是本地的127.0.0.1走的是 loopback 接口没有真实网络层那些限制。验证顺序应该是本地验证QMT 进程、UDP 桥、WebSocket 客户端都在本机确认全链路数据流。跨进程验证在另一个 Python 进程里订阅 WebSocket确认数据能跨进程传递。跨机验证把 WebSocket 客户端放到另一台机器确认防火墙、路由和绑定地址都没有问题。跨机验证时最常碰到的问题有两个UDP 发送端绑定了127.0.0.1导致局域网内的桥接进程收不到。WebSocket Server 绑定了127.0.0.1导致局域网内的下游消费进程连不上。所以设计配置时一定要把绑定地址做成可配置项。本地测试用127.0.0.1生产环境用0.0.0.0或实际网卡地址。4. 真正决定长期稳定性的是这些参数与策略代码能跑通只是第一步。你真正要盯紧的是下面这些容易忽略的细节。4.1 UDP 侧的缓冲、粘包与丢包处理UDP 没有粘包问题因为每个send对应一个数据报文接收端每次recvfrom得到一个完整报文。但有一个问题要和 TCP 区分开UDP 不保证顺序也不保证送达。实际操作里你会遇到接收缓冲区满导致丢包检查系统 UDP 接收缓冲区大小/proc/sys/net/core/rmem_max在 Linux 上通常需要调大。单条行情报文过大QMT 的行情数据如果字段很多JSON 序列化后的字节数可能超过普通 UDP 报文 1500 字节的限制这时容易触发 IP 分片。建议控制单条报文字节数或者拆包。短时间内大量小包小包会消耗大量系统调用次数CPU 占用明显偏高。如果行情频率很高可以考虑在发送端做聚合把多条行情打包成一个数组发送。丢包是 UDP 场景下的常态。桥接端要记录一个单调递增序号缺失的序号代表发生过丢包消费端可以据此判断行情的连续性。4.2 WebSocket 侧的连接管理、心跳与广播策略WebSocket 连接是长连接长连接最容易出现的问题就是“假死”客户端看起来还连着实际上服务端的send已经写不进去了但因为 TCP 缓冲区还在蓄水API 并不直接报错。解决方法在桥接进程侧增加心跳检测定时给客户端发送心跳消息。WebSocket 协议自带 ping/pong 机制建议使用协议级心跳更底层。如果检测到客户端在 N 秒内没有响应直接断开连接释放资源。广播策略也有讲究。简单方案是把每一条 UDP 行情发给所有连接这在小规模场景没问题。但如果你同时给多个客户端发大量行情而某个客户端处理能力不足就会拖慢整个桥。更好的做法是引入“订阅分组”每个客户端连上来时告知自己关注哪些标代码桥接进程只向对应客户端推送它订阅的行情。这相当于在桥里做了一层简单的消息路由。4.3 订阅过滤与消息裁剪QMT 的全市场行情量级非常大而很多时候你并不需要全市场。订阅过滤可以在三个层级做QMT 订阅端只订阅需要的标代码减少进入 UDP 链路的数据量。桥接进程按客户端订阅做过滤减少 WebSocket 推送量。序列化时只保留必要字段去掉冗余的行情明细。订阅过滤做得好可以显著降低 CPU、内存和带宽消耗。但要注意过滤逻辑一旦复杂桥接进程本身会成为瓶颈。如果行情量实在太大更合理的设计是先用 Redis Stream 或消息总线做一次汇聚再让下游按需订阅而不是让桥接进程承担所有路由职责。5. 出问题时按这条链路排查跨进程桥接的问题排查最忌讳的是凭感觉乱改。我建议按下面的链路逐层排查。5.1 阶段一确认行情有没有进入桥先回答一个问题UDP 桥接进程有没有收到数据在这之前先确认 QMT 回调函数本身有没有触发。很多 QMT 用户在桥接链路外排查询情问题时半天找不到原因最后发现是订阅代码没有生效。检查顺序在 QMT 回调函数里打印日志确认回调触发。在 UDP 发送端打印发送计数器确认报文真的发出了。在 UDP 桥接进程里打印接收计数器确认报文真的收到了。在桥接进程里打印解析后的字段确认序列化格式正确。如果第三步没有数据问题出在 UDP 链路。重点看UDP 端口是否绑定在被监听端口。本机用netstat -ulnp能否看到监听进程。防火墙是否拦截了 UDP 端口。5.2 阶段二确认 WebSocket 推送链路UDP 收到了但下游消费者没数据。这时候要去查 WebSocket 链路。检查顺序用websocat或浏览器直接连ws://localhost:13000确认连接建立。在ws_handler里打印连接数确认客户端确实连接上了。打印connected_clients的成员数量确认连接是否被正确注册。检查safe_send是否捕获到了发送异常。这里最容易出错的是连接注册时机。你的客户端可能连上了但因为ws_handler里注册逻辑写在wait_closed()之后导致连接从不加入广播集合数据自然推不出去。5.3 阶段三网络、系统资源与平台边界如果本地全通局域网不通优先检查WebSocket Server 绑定地址是否为0.0.0.0。UDP 发送端是否绑定在127.0.0.1或只监听 loopback 的网口。防火墙是否放通了对应的 UDP 和 WebSocket 端口。是否还有另一层 Docker 或 NAT 转发导致端口映射错误。系统内核对 UDP 缓冲区的大小是否够用。如果发现桥接进程 CPU 高、内存持续上涨优先看是否有 WebSocket 客户端连接失败但未被清理。消息队列是否一直在积压说明消费端处理能力不足。是否每条消息都做了一次 JSON 序列化有没有重复序列化的浪费。6. 适用边界和长期工程化建议最后这部分想做一个冷静的判断这个方案到底适合谁不适合谁以及如果要长期使用还需要补什么。6.1 哪些人最该用这个方案哪些人不适用最该用这个方案的场景你有一个 QMT 行情采集进程已经拿到的行情需要分发给 2 到 5 个消费者。你的消费者里面有 Web 前端需要标准 WebSocket 协议接入。你能接受行情偶尔丢包不需要逐笔完全可靠的对账。你想把行情采集和行情分发解耦让 QMT 进程尽量只做订阅。不建议用这个方案的场景你需要每笔行情都可靠不丢且要能回溯补数据。这时候应该用消息队列把 UDP 换成 Kafka 或 Redis Stream保证持久化和消费确认。你的下游消费者数量很多且各自有完全不同的订阅需求。这时候应该加消息总线或中间件而不是让桥接进程承担复杂路由。你的业务对行情延迟要求在微秒级别。UDP 转发和 WebSocket 推送多出来的网络栈开销可能不适合。6.2 如果要做成生产级服务还差哪几块拼图代码跑通只是开始。如果要长期稳定运行我的建议是按“先跑通、再优化、最后工程化”这个顺序来。工程化阶段需要补齐以下几块日志桥接进程必须有结构化日志至少记录 UDP 收包数、WebSocket 连接数、广播次数、丢弃次数。这样复盘时才不会两眼一抹黑。进程守护桥接进程需要一个守护机制。UDP 桥跑挂了行情分发立刻断掉。用 systemd / supervisor 做自动拉起比手动重启可靠得多。配置外部化UDP 端口、WebSocket 端口、绑定的 IP、订阅过滤规则都要从配置文件或环境变量读取而不是写死在代码里。监控指标至少要有三个指标UDP 收到速率、WebSocket 推送速率、消息队列积压量。这三个指标能帮你定位 80% 的性能问题。合理重连QMT 客户端和 WebSocket 客户端都需要重连机制。尤其是 QMT 端行情中断后必须能自动恢复订阅。还有一个容易被忽略的事情桥接进程的版本管理和接口兼容。当 QMT 行情结构升级消息里新增了字段旧桥接进程能否正确忽略未知字段建议在消息解析时用容错模式新增字段不阻塞旧进程消费。6.3 把“协议转换”当成一次架构抽象来看回到文章开头那句判断。UDP 转 WebSocket 的真正价值不是帮你省掉几行代码也不是把行情从一种格式换成另一种格式。它的价值是让你在 QMT 进程外面开了一层独立的、可扩展的分发面。在这层分发面上你可以继续接前端、接日志系统、接监控面板甚至接另一个策略服务。只要桥接进程还在跑行情就能持续流向下游。协议转换只是手段让实时行情在多进程之间流动才是目的。如果你准备做类似的事建议第一版不要追求复杂功能。先跑通一条最简单链路本地 UDP 发出本地桥接进程收到WebSocket 客户端收到数据。这条链路通了再逐步加过滤、加心跳、加监控、加局域网支持。行情桥这种基础设施永远是先能用再稳定最后才谈优化。