tradingSystem/QMT_SIDE_CONTROL_PATH.md

5.0 KiB
Raw Blame History

给 QMT 侧:ack_seq 不该走控制面内联(阻塞问题的回归)

先说好的部分

server1.py / handlers1.py 里对「接收循环不得 await 业务处理」的修法是对的:

  • CONTROL_TYPES 内联同步处理,ping 能立刻回 pong
  • 业务消息走 asyncio.create_task(_dispatch_business(...))run_in_executor,不占接收循环
  • 任意入站帧都重置 watchdog,不再因业务耗时误判空闲

这一改把「place_order 卡住 → 没人回 pong → 双向 15 秒超时」堵死了,方向完全正确。

ack_seq 放进 CONTROL_TYPES 是个回归

CONTROL_TYPES = frozenset({'ping', 'ack_seq'})

server1.py 对控制面的处理是在事件循环里同步调用

if MessageHandlers.is_control(msg_type):
    self.handlers.handle_control(envelope)   # 同步, 阻塞整个事件循环
    continue

handle_control 里两个分支的性质完全不同:

类型 做的事 耗时
ping publisher.publish('pong', {}) 纯内存,微秒级 —— 适合内联
ack_seq self.store.ack(seq) Redis I/O同步 —— 不适合内联

store.ack(seq) 一阻塞,整个事件循环停转:收不了帧、_sender_loop 也发不出去。 PMS 侧看到的就是「15 秒没收到任何消息」,于是断线重连。

为什么现在才炸

PMS 按协议 §4.5 每 2 秒发一次 ack_seq。在 ack_seq 被实现之前QMT 侧 Redis 里的上行消息一条都没清过,已经积了四千多条。现在每 2 秒调一次 store.ack 第一次清理的量很大,事件循环被按在地上摩擦。

时间线也对得上:

时刻 状态
13:27:59 拒绝 0 · 重连 1 · 时延偏差 +78ms · 连接稳定
~14:00 部署「控制面内联」改动
14:04:33 重连 7 · 偏差 +30213ms · 15 秒静默
14:05:39 重连 8 · 偏差 +97ms · 15 秒静默

那个 +30213ms 不是时钟漂了 30 秒(一分钟后又回到 +97ms真实时钟不会这么跳有一条 pong 在你们的发送队列里躺了 30 秒才发出来 —— 正是事件循环被堵住的证据。

建议改法

ack_seq 没有任何时延要求。协议 §4.5 的原话是「QMT 收到 ack_seq{seq: N}清理」, 晚几秒清理没有任何影响。把它挪回业务面即可:

CONTROL_TYPES = frozenset({'ping'})       # 只留 ping

handle() 里已有 if msg_type in CONTROL_TYPES: return self.handle_control(...) 的兜底, ack_seq 会自然落到业务分支;补一个显式分支更清楚:

if msg_type == 'ack_seq':
    seq = int(payload.get('seq') or 0)
    if seq > 0:
        self.store.ack(seq)
    return

这样 ack_seqrun_in_executor,清 Redis 再慢也不影响 pong

首次清理建议分批:积压四千多条,一次性 ZREMRANGEBYSCORE 之类的操作即使在 线程池里也可能跑很久,分批(比如每次 500 条)更稳。

另一处值得顺手看看

_sender_loop 的队列绑定与 socket 绑定不一致:

sender = asyncio.create_task(self._sender_loop(websocket))   # socket 在创建时绑死

async def _sender_loop(self, websocket):
    while True:
        env = await self._send_queue.get()        # 队列每次循环重新读实例属性
        await websocket.send(...)                 # 但发到的是绑死的那个 socket

self._send_queue 在每次新连接时被重新赋值。旧连接的 sender 任务若尚未被 cancel finally 是异步执行的,close() 返回不代表旧任务已收尾),它下一轮 get() 拿到的 是新队列里的消息,却发往旧的已关闭 socketsend 抛异常 → break 这条消息就此丢失,而新连接那边在等一个永远不来的 pong

重连越频繁越容易撞上,且会自我维持。建议把队列作为参数传进去,与 socket 同生命周期:

queue = asyncio.Queue()
self._send_queue = queue
sender = asyncio.create_task(self._sender_loop(websocket, queue))

async def _sender_loop(self, websocket, queue):
    while True:
        env = await queue.get()
        ...

send_fn 同理,闭包捕获 queue 局部变量而不是读 self._send_queue

一个待确认的口径

handlers1.pypong 仍然是 publish('pong', {}, with_seq=True)。之前提过 「要不要给 pong 免掉 seq」看起来没改这没问题§5 原文就是「所有上行消息带 seq」 带着才是合规的)。只是想确认一下这是有意保留,还是漏了 —— PMS 侧两种都能正确处理, 只影响 pms_qmt_inbox 的行数,不影响正确性。

便于定位的日志

下次断连时,麻烦看一眼这几行有没有出现:

  • 发送失败: ..._sender_loop 里的 break —— 出现就说明发送侧死了)
  • 15s 未收到消息,断开连接(你们的 watchdog
  • PMS 已连接 / PMS 连接清理完成 的先后顺序(错序说明新旧连接在打架)