执行摘要
- 一句话:修复disaggregation receiver ZMQ socket未清理
- 推荐动作:值得精读,体现了分布式系统中socket生命周期管理的设计权衡:通过ZeroMQ选项实现故障快速发现与资源回收。但必须理解
RECONNECT_IVL=-1配合SNDTIMEO的必要性,否则可能导致整个调度循环卡死。建议与后续修复#31025一并阅读。
功能与动机
关联Issue #28596 描述了prefill节点重启时因ZMQ socket未清理导致端口与NCCL冲突。解码端心跳失败处理虽清理了元数据,但未关闭或移除缓存的ZMQ PUSH socket(CommonKVReceiver._socket_cache),造成端口泄露。
实现拆解
- 修改
CommonKVReceiver._connect()方法:在新建ZMQ PUSH socket后增加两个socket.setsockopt调用——zmq.RECONNECT_IVL=-1(禁用ZeroMQ自动重连,避免死节点被无限重试)和zmq.LINGER=0(关闭时立即返回,不等待未发送消息)。
- 范围精简:应维护者要求移除了原方案中的
SNDTIMEO、TCP_KEEPALIVE等额外选项,仅保留上述两个核心配置。
- 保持现有锁机制:
disconnect_endpoint方法继续使用每端点锁保护sock.close(),未采纳移除锁的建议。
- 移除单元测试:维护者认为无需独立测试,测试文件被回滚。
关键文件:
python/sglang/srt/disaggregation/common/conn.py(模块 网络层;类别 source;类型 core-logic;符号 _connect): 核心变更文件,在_connect方法中新增两行socket选项配置,直接影响disaggregation场景下ZMP socket的故障清理行为。
关键符号:CommonKVReceiver._connect
关键源码片段
python/sglang/srt/disaggregation/common/conn.py
核心变更文件,在_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]
评论区精华
风险与影响
- 风险:主要风险来自
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 保护, 线程竞争风险(锁机制)
关联脉络
参与讨论