执行摘要
- 一句话:为 NIXL 后端接入延迟解码端 KV 释放的 ABORT_ACK 机制
- 推荐动作:值得精读,尤其适合负责 PD / KV 传输正确性的工程师。三个设计决策值得学习:① 用「出队时计数」构造静默判定(
Failed and outstanding == 0 才可安全 ack)的并发推理;② decode 侧为什么必须用独立 ZMQ 线程而非 NIXL notif 消费 ack(num_threads=0 下 notif 没有独立进度线程);③ 异常路径故意不 ack、交给超时兜底的「宁可保守不冒险」取舍。建议先在小规模 NIXL 部署做一次 abort-mid-transfer 端到端验证,再在生产开启该特性。
功能与动机
PR body 明确指出两个动机:一是「latent footgun」——decode 侧与后端无关,启用 SGLANG_DISAGGREGATION_DEFERRED_DECODE_KV_RELEASE 后 NIXL 部署 already holds pages,但 NIXL 从不发送 ABORT_ACK 也没有接收路径,is_abort_release_safe 永远不满足,每次 abort 都等满超时——安全但完全退化(原话:'safe (holding is the safe direction), but fully degraded, pinning KV and back-pressuring new requests under abort load');二是 NIXL 的传输天然异步(agent.transfer(handle) 发出后由 worker 轮询 check_xfer_state(handle)),in-flight 窗口真实存在,危险性比 mooncake 默认的同步 batch_transfer_sync 更甚;nixl/conn.py::_handle_abort_notification 原本就带着 TODO。上游 #35049 只给 mooncake 接上了 prefill ack。
实现拆解
变更入口是 python/sglang/srt/disaggregation/common/conn.py 与 python/sglang/srt/disaggregation/nixl/conn.py,配套改动 mooncake/conn.py 保持行为一致,测试集中在两个 NIXL 测试文件。拆解如下:
-
共享化 ack 机制(common/conn.py,+44 行):把 _prefill_unique_rank、_send_abort_ack、_maybe_ack_drained_abort 三个方法及 _deferred_ack_targets 映射从 MooncakeKVManager 上提为 CommonKVManager 的成员,并新增 register_deferred_ack_target(),其 docstring 明确「先 mark Failed 再登记」的顺序契约——若在 room 仍接受 chunk 时登记,worker 可能先 ack,而后续 chunk 又把数据写进 decode 已释放的页。mooncake 删除本地副本改为继承,行为不变;CommonKVSender.clear() 已有的 target 清理逻辑让 NIXL 免费获得同样清理。
-
NIXL prefill 侧接入(nixl/conn.py,+52/-15 行):_handle_abort_notification 解析 decode 在 _send_abort_notification 中一直携带的 ip/port(NIXL 此前忽略),对 active room 在状态翻转(置为 Failed)之后登记 ack target 并立即尝试一次 ack;对已 concluded 的 room 立即 ack。bootstrap 线程绝不 ack 仍在活跃的 room(chunk 已出队但未计数时会与 ack 竞态)。transfer_worker 在两个时点补发 ack:跳过点(room 已 Failed,本 chunk 不会写入任何数据)与既有 while handles: check_xfer_state == DONE 屏障之后(写入已落定);异常路径故意不 ack——DONE 屏障在首个 ERR 上抛出,兄弟 handle 可能仍在写,room 回落到 decode 侧超时兜底。
-
修复两处 ack 滞留窗口(共享逻辑,两个后端均生效):根因是 outstanding == 0 不是静默判定——chunk 在状态检查之后才计数,导致「chunk 已出队、即将写入」时计数器仍为 0,bootstrap 线程因此不能安全地发 ack。现在改为出队时、状态检查前计数,chunk 要么已计数并写入、要么被跳过不计数,二者必居其一;Failed and outstanding == 0 才真正意味着不再有写操作。abort 路径在登记后立即尝试一次 ack,覆盖「worker 先于登记到达跳过点」与「chunk 全部排空后 worker 不再回访该 room」两个窗口。计数提前对既有 poll() Success 持有与 teardown 条件更保守(增/减仍配对),mooncake 的 transfer_worker 同步应用同一计数顺序。
-
decode 侧 listener 泛化与测试配套(nixl/conn.py + 两个测试文件):_start_decode_staging_thread 改为 _start_decode_listener_thread,同一 ZMQ 线程同时消费 STAGING_REQ 与 ABORT_ACK,ack 按 prefill rank 聚合进 note_abort_ack;启动条件从 staging 扩展为 staging 或 deferred release 任一开启。必须用独立线程而非 NIXL notif:decode agent 以 num_threads=0 运行,notif 只在活跃 receiver 的 poll() 内被消费,decode 空闲(conclude_state 已置位)时会漏。测试方面新增 test/registered/unit/disaggregation/test_nixl_deferred_kv_release.py(12 个用例,用 object.__new__ 构造裸 manager、以 lambda 捕获 ack),并把 test_nixl_backend_basic.py 两组 fixture 显式置 enable_deferred_decode_kv_release = False 锁定 legacy 无 ack 行为;mooncake 侧套件一并纳入回归(共享 worker 逻辑变更)。
关键文件:
python/sglang/srt/disaggregation/nixl/conn.py(模块 传输层;类别 source;类型 core-logic;符号 _start_decode_listener_thread, decode_listener_thread, _handle_abort_notification, transfer_worker): NIXL 后端的核心接入点:_handle_abort_notification 解析 abort 通知并注册/补发 ack,transfer_worker 在出队计数、跳过点与 DONE 屏障后补发 ack,decode 侧把 staging 专用线程泛化为同时消费 ABORT_ACK 的 listener;也是计数顺序修复与两处 ack 滞留窗口修复的直接载体。
python/sglang/srt/disaggregation/common/conn.py(模块 公共层;类别 source;类型 core-logic;符号 _prefill_unique_rank, _send_abort_ack, _maybe_ack_drained_abort, register_deferred_ack_target): 共享 ack 机制上提至此:_prefill_unique_rank、_send_abort_ack、_maybe_ack_drained_abort 与 _deferred_ack_targets、register_deferred_ack_target 成为 CommonKVManager 的一部分,是 mooncake 与 nixl 两个后端共同依赖的正确性契约所在,也定义了 mark-Failed-first 的顺序要求。
python/sglang/srt/disaggregation/mooncake/conn.py(模块 传输层;类别 source;类型 dependency-wiring;符号 transfer_worker): 行为保持不变的依赖重接:删除本地重复实现(-48 行)改为继承共享机制,并把 transfer_worker 的 chunk 计数顺序同步改为「出队时计数」,以消除与 NIXL 相同的 ack 滞留窗口;是确认共享改动不破坏 mooncake 语义的关键对照文件。
test/registered/unit/disaggregation/test_nixl_deferred_kv_release.py(模块 单元测试;类别 test;类型 test-coverage;符号 TestDeferredAckTargets, TestNixlAbortNotification, test_ack_held_until_outstanding_drains, test_ack_fires_at_most_once): 新增 12 个用例覆盖完整 ack 状态机:ack 持有直到 outstanding 排空、至多发送一次、未登记 room 为 no-op、rank 公式与 Success 同步一致、abort 通知各分支(active / concluded / feature-off / legacy 2-frame abort / 非 ABORT 透传)、按 rank 聚合;是验证并发不变式的主要手段。
test/registered/unit/disaggregation/test_nixl_backend_basic.py(模块 单元测试;类别 test;类型 test-coverage): 既有 NIXL fixture 通过 object.__new__ 构造 manager,需显式设置 enable_deferred_decode_kv_release = False 以锁定 legacy 无 ack 行为,防止新特性在旧用例中静默改变断言语义。
关键符号:_start_decode_listener_thread, decode_listener_thread, _handle_abort_notification, transfer_worker, _maybe_ack_drained_abort, _send_abort_ack, register_deferred_ack_target, _prefill_unique_rank, note_abort_ack
关键源码片段
python/sglang/srt/disaggregation/nixl/conn.py
NIXL 后端的核心接入点:_handle_abort_notification 解析 abort 通知并注册/补发 ack,transfer_worker 在出队计数、跳过点与 DONE 屏障后补发 ack,decode 侧把 staging 专用线程泛化为同时消费 ABORT_ACK 的 listener;也是计数顺序修复与两处 ack 滞留窗口修复的直接载体。
# python/sglang/srt/disaggregation/nixl/conn.py
# decode 侧:把只处理 STAGING_REQ 的 ZMQ 线程泛化为 listener,
# 同时消费 ABORT_ACK。必须用独立线程而不是 NIXL notif —— decode agent
# 以 num_threads=0 运行,notif 只能在活跃 receiver 的 poll() 内被消费,
# decode 空闲(conclude_state 已置位)时会漏掉 ack。
def _start_decode_listener_thread(self):
"""Decode-side ZMQ listener for STAGING_REQ and ABORT_ACK."""
def decode_listener_thread():
while True:
msg = self.server_socket.recv_multipart()
if msg[0] == b"STAGING_REQ":
self._handle_staging_req(msg)
continue
if msg[0] == b"ABORT_ACK":
# 排空 ack:消息格式为 ABORT_ACK + room + prefill_rank;
# note_abort_ack 按 room 聚合,room 未在持有时直接丢弃,
# 因此迟到的 ack 不会污染已被复用的 bootstrap_room。
if len(msg) >= 3:
self.note_abort_ack(
int(msg[1].decode("ascii")),
int(msg[2].decode("ascii")),
)
continue
logger.warning(
"decode_listener_thread: unexpected message tag %s",
msg[0][:20],
)
threading.Thread(target=decode_listener_thread, daemon=True).start()
评论区精华
唯一的实质讨论线程由 YAMY1234 提出,指向 nixl/conn.py:2692:
YAMY1234:Could this ACK too early if the room has been cleared while NIXL transfers are still in flight? The ERR path says sibling handles may still be writing, so would a later ABORT hit the unknown-room branch and ACK immediately?
ShangmingCai:Good catch — this is a real hole, and I've fixed it in 30a50a2. …… The DONE barrier raises on the first ERR handle, so sibling handles are abandoned and may still be writing — and because it raises, the self._staging_outstanding[room] -= 1 below it never runs. …… the sender then concludes Failed → failure_exception() → clear() → which pops the room from request_status. A later ABORT for that room therefore hits the concluded/unknown branch and acked immediately……
ShangmingCai 承认是真实漏洞并在 30a50a2 修复,修复原则即提交信息所述:don't ack a cleared room that still has work counted——unknown 不隐含静默。两位 reviewer(YAMY1234:'LGTM overall — one question inline.';ishandhanani:'Overall LGTM')均 APPROVED。
- room 已清理但传输仍在途时可能过早 ACK (correctness): ShangmingCai 确认是真实漏洞并在 30a50a2 修复:DONE 屏障在首个 ERR handle 上抛异常,
_staging_outstanding[room] -= 1 不会执行;except 块只置 Failed,sender 随后 conclude Failed → failure_exception() → clear() 把 room 从 request_status 弹出;后续 ABORT 命中 concluded/unknown 分支会立即 ACK。修复原则如提交信息所述:don't ack a cleared room that still has work counted——unknown 不隐含静默。
风险与影响
- 风险:
- 并发正确性敏感:abort/ack 状态机跨 bootstrap 线程与多个 transfer worker,修复建立在「出队计数先于状态检查」这一不变式上;任何未来重构若改回计数顺序,会静默重开两处 ack 滞留窗口。多节点 / TP>1 dummy-rank 扇出未在真实并发下验证(PR body 自述)。
- 缺少 NIXL 端到端验证:状态机与接线仅有单元测试;abort-mid-transfer 在真实 NIXL 部署上的 fast-ack-vs-timeout 比值未测量,与 #35049 的 gap 相同。
- 共享逻辑扩大影响面:
common/conn.py 的改动同时作用于 mooncake,计数提前使 poll() Success 持有与 teardown 条件更保守,但任何计数错配都会同时影响两个后端——mooncake 侧套件已纳入回归验证。
- 异常路径依赖超时兜底:DONE 屏障 ERR 后
_staging_outstanding 不递减、room 不再被 worker 回访,只能等 decode 侧超时释放;若超时配置偏小,仍存在「兄弟 handle 写入已释放页」的理论窗口(#35049 既有设计权衡,本 PR 未改变)。
- 默认关闭降低回归面:feature 关闭时两个后端行为与之前完全一致,未被显式开启的部署不受影响。
- 影响:
- 用户/部署方:启用
SGLANG_DISAGGREGATION_DEFERRED_DECODE_KV_RELEASE 的 NIXL PD 部署从「每次 abort 空等满超时、钉住 KV 并反压新请求」变为「传输排空即释放、按 rank 聚合 ack 后立即回收」,abort 场景下的 KV 占用与队列阻塞显著改善;mooncake 部署仅需升级,行为不变。
- 系统:改动位于 disaggregation 连接层核心(
CommonKVManager 与 NIXL / Mooncake 两个后端),触及 abort 处理与 KV 页释放的时序正确性,属于高危路径。
- 团队:新增 12 个单元测试覆盖 ack 状态机全部分支;两个后端共享 ack 机制后,后续维护(如新增传输后端)必须同时满足 mark-Failed-first 与 dequeue-counting 两个不变式。
- 风险标记:核心路径变更, 并发正确性敏感, 缺少 NIXL 端到端验证, 多节点/TP>1 未验证
关联脉络
- PR #35049 [PD] Deferred decode-side KV release: 本 PR 的直接上游:引入 env-gated 延迟 decode 侧 KV 释放机制但只给 mooncake 接通 prefill ack;本 PR 将其扩展到 NIXL 并修复其遗留的两处 ack 滞留窗口。
- PR #35424 [Fix] Scale the req_to_token row headroom by attn_dcp_size: 同属 disaggregation / KV 页安全主题:按 DCP 规模缩放 req_to_token 行头空间修复越界,均在加固 KV 页复用时序与边界正确性。
- PR #35298 [Fix] DCP: advertise the logical KV-event block size: 同属 disaggregation KV 传输正确性修复系列(修复 DCP 下事件块大小广告错误导致的路由命中率归零),与本 PR 处于同一功能域。
参与讨论