Prhub

#29570 Fix disaggregation receiver ZMQ cleanup

原始 PR 作者 ronhuafeng 合并时间 2026-06-29 16:22 文件变更 1 提交数 2 评论 17 代码增减 +2 / -0

执行摘要

修复 disaggregation receiver ZMQ socket 未清理

关联Issue #28596 描述了prefill节点重启时因ZMQ socket未清理导致端口与NCCL冲突。解码端心跳失败处理虽清理了元数据,但未关闭或移除缓存的ZMQ PUSH socket(CommonKVReceiver._socket_cache),造成端口泄露。

值得精读,体现了分布式系统中socket生命周期管理的设计权衡:通过ZeroMQ选项实现故障快速发现与资源回收。但必须理解RECONNECT_IVL=-1配合SNDTIMEO的必要性,否则可能导致整个调度循环卡死。建议与后续修复#31025一并阅读。

讨论亮点
  • 锁与线程安全(design):gemini-code-assist建议在disconnect_endpoint中不加锁直接关闭socket,但作者基于PyZMQ文档指出Socket.close()非线程安全,坚持保留每端点锁以避免与send_multipart并发时的竞争。达成共识,维持原锁。
  • 变更范围精简(design):维护者ShangmingCai要求只保留RECONNECT_IVLLINGER,移除SNDTIMEO及TCP keepalive选项。作者接受并精简。
  • 单元测试必要性(testing):ShangmingCai认为无需新增测试,作者删除测试文件。

实现拆解

  1. 修改CommonKVReceiver._connect()方法:在新建ZMQ PUSH socket后增加两个socket.setsockopt调用——zmq.RECONNECT_IVL=-1(禁用ZeroMQ自动重连,避免死节点被无限重试)和zmq.LINGER=0(关闭时立即返回,不等待未发送消息)。
  2. 范围精简:应维护者要求移除了原方案中的SNDTIMEOTCP_KEEPALIVE等额外选项,仅保留上述两个核心配置。
  3. 保持现有锁机制disconnect_endpoint方法继续使用每端点锁保护sock.close(),未采纳移除锁的建议。
  4. 移除单元测试:维护者认为无需独立测试,测试文件被回滚。
文件 模块 状态 重要度
python/sglang/srt/disaggregation/common/conn.py 网络层 modified 5.12

关键符号

CommonKVReceiver._connect

关键源码片段

python/sglang/srt/disaggregation/common/conn.py core-logic

核心变更文件,在 `_connect` 方法中新增两行 socket 选项配置,直接影响 disaggregation 场景下 ZMP socket 的故障清理行为。

@classmethod
def _connect(cls, endpoint: str, is_ipv6: bool = False):
    with cls._global_lock:
        if endpoint not in cls._socket_cache:
            sock = cls._ctx.socket(zmq.PUSH)
            if is_ipv6:
                sock.setsockopt(zmq.IPV6, 1)
            # 禁用 ZeroMQ 自动重连,避免死节点占用缓存 socket
            sock.setsockopt(zmq.RECONNECT_IVL, -1)
            # 设置 LINGER=0,确保 close() 立即返回,不等待未发送消息
            sock.setsockopt(zmq.LINGER, 0)
            sock.connect(endpoint)
            cls._socket_cache[endpoint] = sock
            cls._socket_locks[endpoint] = threading.Lock()
        return cls._socket_cache[endpoint], cls._socket_locks[endpoint]

评论区精华

disconnect_endpoint 中关闭 socket 的锁使用 设计

gemini-code-assist 认为 sock.close() 在 PyZMQ 中线程安全,建议不加锁关闭以避免心跳线程阻塞;ronhuafeng 引用 PyZMQ 文档指出 Socket.close() 非线程安全,且与 send_multipart 存在竞争,必须保留锁。

结论:维持现有锁保护,不采纳无锁建议。 · 已解决

移除多余的 socket 选项只保留核心两个 设计

ShangmingCai 要求只保留 RECONNECT_IVL 和 LINGER,移除 SNDTIMEO、TCP_KEEPALIVE 等冗余选项,以减少变更影响。

结论:作者更新代码,仅保留两行配置。 · 已解决

添加单元测试是否有必要 测试

ShangmingCai 认为无需新增单元测试,作者随即删除了测试文件。

结论:测试被移除,不引入。 · 已解决

风险与影响

主要风险来自RECONNECT_IVL=-1:当对端PULL socket死亡后,PUSH socket进入mute状态,后续send_multipart若无SNDTIMEO会无限阻塞(此PR未设置SNDTIMEO)。该问题在后续PR #31025中修复。此外,LINGER=0可能导致未发送消息丢弃,但影响有限。锁保护虽避免竞争,但若send_multipart长时间阻塞,心跳检测线程也可能因等待锁而延迟。

影响范围:仅限disaggregation部署模式下decode侧的receiver socket创建逻辑。对正常请求路径无影响(复用缓存socket),仅在故障恢复时改变socket清理行为。需注意用户如果自定义ZMQ选项可能被覆盖。

RECONNECT_IVL=-1 缺 SNDTIMEO 保护 线程竞争风险(锁机制)

关联 Issue

#28596 [Bug] [Disaggregation] ZMQ sockets are not cleaned up after prefill node failure, causing port conflicts with NCCL on prefill restart

完整报告

参与讨论