Completion Semantics of NCCL Failure Recovery
After the sender sees completion, the receiver may still not have consumed the message. Should failure recovery trust the sender's completion, or the application effect the receiver has already produced?
This note is about one concrete question: when an NCCL transfer fails halfway, how do we decide whether a message must be sent again?
Suppose the sender has emitted a piece of KV cache or an expert output, and the CUDA event for that send has already completed. The receiver process is then killed. Recovery has two choices, resend or do not resend. Not resending is correct only if sender completion already proves that the receiver finished the work. Resending is correct only if the receiver has not yet used the message. Both premises can fail at once.
The conjecture we started from was simple. The completion NCCL exposes to the sender, together with the communicator’s error state, might already be enough to define the retry boundary. If that were true, the layer above would not need per-message state. Recovery would resend only the operations whose completion the sender had not yet observed.
What has to be tested is therefore not whether NCCL can report an error. It is whether this equivalence holds:
sender completion ≈ the receiver has already produced the application effect
The rest of the experiment is organized around that relation.
Why this starts to matter
When large-scale training hit a communication fault, the usual response was to fail the whole group and restart from a checkpoint. That is expensive, and the semantics are simple: once the step is discarded, nobody has to know how far any in-flight message got.
Current systems are less willing to pay that cost. torchft can drop a failed replica at a step boundary. NCCL’s ncclCommShrink and ncclCommRevoke let a communicator be rebuilt after a fault. Inference also keeps long-lived state. Expert-parallel token dispatch, KV-cache transfer between prefill and decode, and long RL rollouts all want to keep work that already finished, instead of wiping the step and starting over (Qin et al., FAST 2025).
flowchart LR
A["Then: rank fails"] --> B["Drop the step"] --> C["Restart from checkpoint"]
D["Now: rank fails"] --> E["Revoke / rebuild"] --> F["Keep finished work"]
F --> G["Which messages took effect?"]
That picture is where the problem comes from. Message-level completion does not matter while the whole step is thrown away. The moment the system tries to keep partial progress, message-level semantics become the basis of correct recovery. A KV block that a decode worker has already attached to later token generation can create duplicate state if it is replayed. The opposite miss is just as bad. Data that only reached an intermediate buffer, and was never consumed, is dropped for good if sender completion is treated as a reason not to resend.
Completion, delivery, and consumption are not the same event
Three events have to stay separate. The first is local completion. The sender has handed the send to the NCCL and CUDA pipeline, and the send buffer may now be reused. ncclSend itself is asynchronous. It only enqueues work on a CUDA stream. The caller usually records a CUDA event afterward, so the event means “the send-side work has finished,” not “the peer application has processed the bytes.” The second is delivery. The bytes have landed in receiver memory. The third is consumption. A receiver kernel or the application logic has read those bytes and turned them into an effect that cannot be ignored, such as attaching a KV-cache block to later decode, or adding an expert output into the final result.
A parcel is a useful picture. Local completion is the drop box accepting the package, so the sender can leave. Delivery is the package arriving at the address. Consumption is the recipient opening it and using what was inside. In a healthy run the three stay in order, and it is easy to call all of them “done.” Recovery is hard because a fault can cut exactly between them.
flowchart LR
S["1. ncclSend"] --> C["2. Sender completion"] --> D["3. Delivery"] --> U["4. Consumption"]
The object of this note is the layer at which failure recovery places its commit point. A commit at sender completion can believe too early that a message is done. A graph-level or batch-level completion can confirm too late, and replay a message the receiver already consumed.
On an RDMA reliable connection, a completed write means the remote NIC has acknowledged that the data was written into remote memory. That is delivery. It does not include application consumption (InfiniBand specification). After a queue pair enters the error state, outstanding requests return as flush errors, but some of those writes may already have happened on the peer. A failed request does not imply that nothing happened remotely.
This experiment has no GPUDirect RDMA, so the path is longer. The NIC cannot write remote data straight into GPU HBM. The bytes land first in a CPU-visible host bounce buffer, NCCL’s proxy and copy path moves them into the GPU, and only then do the receive kernel and the later consumer kernel use them. That is an extra handoff between “the parcel reached the building” and “someone opened it in the room.” Each extra stage is another window a fault can enter. The gap is a property of this data path, not an abstract argument about words.
The recovery primitives we already have are group-scoped. ncclCommAbort drops every in-flight operation. ncclCommRevoke stops in-flight operations and marks the communicator in progress, after which it can be destroyed or shrunk once it is quiescent. ncclCommGetAsyncError returns one error code for the whole communicator (NCCL communicator API). PyTorch’s ProcessGroupNCCL adds a watchdog that fails every outstanding Work on the group. MPI’s fault-tolerance extension ULFM provides revoke, shrink, and agree. Agree lets surviving processes reach a fault-tolerant consensus on an integer (Bland et al., 2013), but a failed request itself only reports process failure or revocation. UCX’s peer error-handling mode guarantees that every send request ends in success or error, and an error ending does not say whether the peer received the data (UCX).
The window is large in absolute time. In this setup, abort or revoke itself takes about 0.5 s, mostly from the NCCL proxy thread’s 500 ms poll. A 16 B send/recv round trip is only tens of microseconds. The control interval from “one side has decided to give up” to “the other side is actually quiescent” is about four orders of magnitude longer than one message lifetime. The send window allows at most 32 operations in flight, and those 32 can keep crossing the boundaries of sent, delivered, and consumed during that half second. The network did not suddenly become slow. The recovery control plane is coarse and slow, while the data plane is still moving forward.
TL;DR
A communication library telling the sender “I am done” does not mean the receiver application has used the message. On two V100s, EDR InfiniBand, and NCCL 2.32.3 without GPUDirect RDMA, we injected 634 faults. When resend decisions used only sender completion, eager and revoke recovery missed about 1–4 messages per fault that were “send complete” but not yet consumed. A 16-op CUDA Graph, because its completion is that coarse, duplicated about 4–9. After recovery followed the receiver’s actual consumption sequence, duplicates and losses both fell to 0, and the extra cost stayed inside measurement noise. The stable commit point is the application effect, not transport or runtime completion.
What we actually conjectured
The simplest conjecture is that sender completion can stand in for a message’s recovery commit point. In-flight messages then fall into two classes. Those whose completion the sender has seen are not resent. Those still incomplete are sent again after recovery.
The conjecture looks reasonable because, on a healthy run, completion and consumption sit next to each other. NCCL also hides a lot of RDMA, proxy, buffer, and stream detail from the application. A CUDA event is the most natural signal the upper layer can get that “this communication is over.”
It has two openings. First, completion can be too early. The sender believes the send is finished, but the receiver has only received the bytes and has not run the consumer kernel. Second, completion can be too coarse. A CUDA Graph folds many sends into one graph-level event. The receiver may already have consumed part of the batch, while the sender can still only see that the whole graph is not done.
We therefore also ran a stronger control. Every message carries a sequence number. The receiver writes that number into a log only after it has actually consumed the message. On recovery the sender does not consult its own completion. It asks the receiver which sequence numbers have produced an application effect, and it resends only the missing ones. That is the end-to-end argument: the communication layer can report how much was transferred, and only the endpoint knows whether the operation happened in the application’s semantics (Saltzer, Reed, and Clark, 1984).
The experiment compares two commit points:
flowchart LR
A["Fault"] --> B{"Trust?"}
B --> C["Sender completion"]
B --> D["Receiver log"]
C --> C1["done: no resend"]
C --> C2["open: resend"]
D --> D1["consumed: no resend"]
D --> D2["missing: resend"]
C1 --> E["may under-send"]
C2 --> F["batch may over-send"]
D1 --> G["application effect"]
D2 --> G
If sender completion and receiver consumption agree, the two paths should produce the same result. If they systematically diverge inside the fault window, the error shows up directly.
How the experiment separates the three events
Figure 1 is one message lifetime. A kernel on the sender GPU fills the message and ncclSend emits it. The CUDA event after that send is the completion the sender sees. The data is written through an InfiniBand reliable connection into the receiver’s host buffer, then copied to the GPU. After ncclRecv completes, the receiver’s consumer kernel checks and applies the message, then appends the sequence number to a durable log. We define “consumed” by that log. The shaded region is the window where the sender has reported completion and the receiver has not consumed. A fault there, combined with “resend only what is not complete,” loses the message. The note under the figure is the graph path: 16 sends are captured in one CUDA graph, and the sender sees a single completion event for the whole graph.
The hardware is two nodes, one V100-PCIe-32GB each, on EDR InfiniBand, without GPUDirect RDMA, so the data path uses the host bounce buffer. The software is NCCL 2.32.3 and CUDA 12.4. The first 16 B of each message are a 64-bit sequence number and a checksum. Each payload word is determined by the sequence number, and the receiver checks every message. Message size alternates between 16 B and 64 KB across trials. The send window is 32. On the graph path each graph holds 16 messages, with two graphs in flight. A separate TCP control connection carries the communicator unique id, the fault plan, the abort notice, and the log query after the fact.
There are three recovery paths. Eager is per-message send/recv, then ncclCommAbort and a new communicator. The graph path also uses abort. Revoke uses a nonblocking communicator, then ncclCommRevoke, destroy, and a new communicator. There are three faults. F1 sends SIGKILL to the receiver in the middle of the stream, and a supervisor restarts it and brings it back. F2 is the sender aborting itself mid-stream. F3 is SIGSTOP on the receiver for 100 to 800 ms, then SIGCONT. The pause includes NCCL’s proxy thread. Injection time is uniform over the length of a fault-free stream. Detection uses only signals that already exist: control-socket disconnect, ncclCommGetAsyncError, a 500 ms stall, and F2’s self-abort.
We compare two recovery rules. The first uses only the completion the sender can see. It resends every message whose completion the sender has not observed, and the receiver does not deduplicate. The second uses endpoint consumption. Messages carry sequence numbers. The receiver keeps a dedup bitmap and writes the log only after real consumption. After a restart it restores that log. Once communication is back, the sender queries the log and resends only missing sequence numbers. A duplicate is the same sequence number consumed more than once. A loss is a sequence number still absent from the log after recovery. Stale is the count, at the fault, of messages whose sender completion and receiver consumption disagree. Each cell is 210 injections. In F3, a pause shorter than 500 ms trips no detector and only adds latency, so the trigger count is below 210. Fault-free controls, 40 runs each, have all four counters at 0, which rules out harness false positives. A fourth fault, link jitter, would require changing the NIC port state and was not allowed on the shared nodes, so it was not measured.
Where sender completion and receiver consumption diverge
Figure 2 has one panel per fault. The horizontal axis is the recovery path and policy. The vertical axis is the average number of duplicate messages (solid) or lost messages (hatched) per fault that actually triggered recovery. The absolute height matters less than the direction of the error. Eager and revoke lean toward loss. The graph path leans toward duplication. Those two directions match two completion granularities. The first decides too early that a single send is done. The second waits until a whole batch is done.
On the eager path with fixed retry, killing the receiver (F1) loses 4.0 messages per trigger (675 messages across 169 triggers). A pause (F3) loses about 1.0. A sender abort (F2) loses 28 messages and duplicates 1 across 171 triggers. Against a window of at most 32 in-flight operations, 4 is not the whole window being wrong. It is the last few messages caught in the gap where the sender has already released the operation and the receiver has not yet consumed it. That magnitude matches the data path. Most messages are either not yet complete or already consumed. Only a short tail of the pipeline sits in the ambiguous region.
The lost messages share one property. The sender has already seen completion, so recovery will not resend them. The receiver may only have the bytes in a host or GPU buffer, with the consumer kernel not yet run, when the process is killed. Or abort invalidates data that was delivered and not yet consumed. The eager error is not that retry is too aggressive. It is that send completion is too early to use as the acknowledgement point. Revoke has the same shape. F1 loses 2.1 messages per trigger. F2 and F3 accumulate 27 and 91 losses, with no duplicates. Revoke shortens the kill path from eager’s roughly 880 ms to 385–520 ms. That changes how long it takes to stop. It does not change which event counts as the message having happened.
The graph path runs the other way, and the dominant error is duplication. Under F2, 210 triggers produce 1901 duplicate consumptions, about 9 per fault. F1 and F3 are about 3.9 and 6.2. The direct picture is that a CUDA Graph gives the sender one receipt for 16 messages. At the fault, the receiver may already have consumed the first 9, while the graph completion event is still dark. Recovery can only treat all 16 as unfinished and send them again, so the first 9 are applied twice.
Nine is not noise. It is the natural size of a fault in the middle of a 16-op batch. The closer the fault sits to the end of the batch, the more messages have been consumed without a batch completion. Once the graph coarsens completion, the error flips from sending too little to sending too much. If consumption is not idempotent, for example accumulating a gradient, merging an expert output, or advancing a state machine, that duplicate is not just wasted work. It can become a silent numerical error.
The end-to-end column is 0 on all three faults. On the eager path, F1, F2, and F3 trigger 178, 177, and 68 times (F2 has 4 extra injections, 214 in total). Duplicates and losses are both 0. The protocol does not use a more elaborate communication primitive. The receiver records which sequence numbers it actually consumed, and recovery fills the gaps in that ledger. Sender completion is a record of send-side progress. The receiver log is a record of application effect. Moving the recovery basis from the first to the second removes both of the opposite errors above.
The cost is small. On a fault-free stream, turning dedup on changes 16 B throughput by +1.7% and 64 KB throughput by -0.8%, both inside run-to-run noise. Sequence comparisons and bitmap lookups are light per-message metadata next to the data movement and GPU kernel launch. On the graph path we only ran fixed retry online. Offline, adding receiver dedup removes the duplicates. The remaining gaps, 121 on F1 and 6 on F3, are exactly sequence numbers absent from the receiver log, which is the set a log-based resend would identify. That online graph recovery path was not measured separately.
Three failure modes the runs exposed directly
ncclCommGetAsyncError returned success on every triggered fault, including after the peer process had been killed with SIGKILL. The moment a peer is dead and the moment the communicator has a queryable asynchronous error are not the same moment. Detection in the experiment came from the control socket closing, from a 500 ms stall, or from the sender aborting itself. An upper layer that treats this API as an immediate reset, in the sense of a TCP connection reset, waits forever in this environment.
Abort also does not retract kernels already queued on the GPU stream. The receiver still schedules a large number of consumer kernels. Their recv slots were never filled with new data, and each cell shows about 2000–6400 of these empty reads. The communication control plane has stopped, and the GPU execution queue does not rewind itself. We clear the 16 B header before every recv, so those kernels see an empty header and refuse to consume. Another 14 messages had a header that looked new and a payload that was half new and half old. The checksum caught the tear. Without an explicit empty marker and an integrity check, both cases look like legal input.
Cleanup order decides whether recovery can finish. While a CUDA graph exec still holds the communicator, ncclCommAbort does not return for more than 35 s. Destroying the graph exec first, then aborting, returns normally. The lifetime dependency is inverted. The obvious order is to stop communication and then clean up the graph, but the graph still holds the communication resource, so abort waits on an object that will not be released first. For a graph-based inference engine this is more dangerous than a slow path. The wrong cleanup order turns a recoverable fault into a hang inside recovery itself.
How endpoint consumption removes the ambiguity
The rightmost group in Figure 2 draws a sharp boundary. The completion and communicator error that NCCL exposes are not enough to decide per-message recovery. The receiver’s consumption state is enough. Across 634 injections, sequence numbers, dedup, and resend-from-the-receiver-log brought both duplicates and losses to 0. On fault-free streams, 16 B throughput moved by +1.7% and 64 KB by -0.8%, both inside run-to-run noise.
The information that does the work is not a finer NIC completion, and it is not a faster communicator teardown. It is “which sequence number the application has consumed.” The two sides keep different books. Sender completion is an execution log: which operations I handed off. The receiver consumption log is an effect log: which operations actually changed application state. Recovery needs the second book. Reading only the first under-sends on the eager path and over-sends on the graph path. Filling the difference from the second book removes both errors. Under partial failure, the correct recovery boundary crosses the communication library and lands on the end-to-end application effect.
The result also marks what is still open. This is two-rank point-to-point communication. Partial completion of a collective has no per-message counterpart here. There is no GPUDirect RDMA, and the host bounce buffer widens the gap between local completion and consumption. The consumer kernel’s dedup check and log append sit in the same kernel. We did not inject a fault between them, so consumption that has an external side effect still needs its own atomicity story. Link jitter and MSCCL++ were not measured. MSCCL++ 0.9.0 has no abort, revoke, or per-operation status. The source shows a coarser recovery boundary, and that reading is not a substitute for fault-injection data.
What these data support directly
On NCCL 2.32.3 point-to-point send/recv, a completed sender CUDA event means only that the local send has finished. When the receiver dies or is aborted, each fault leaves 1 to 4 messages that were reported complete and never consumed. That is two V100s, EDR InfiniBand without GPUDirect RDMA, a window of 32, and 16 B or 64 KB messages. Deciding what to resend from sender completion drops those messages.
Communication captured in a CUDA Graph reports completion for the whole graph. At the fault the receiver may already have consumed part of the batch. Resending every operation that has not completed duplicates about 4 to 9 messages per fault. One graph here is 16 operations. A coarser grain means more duplicates. Graph communication needs a sequence number on each message. It cannot rely on completion of the graph.
When the peer is killed or paused, or when the local rank aborts itself, ncclCommGetAsyncError keeps returning success. Fault detection needs a separate control channel, or a timeout in the application.
After abort, receive kernels already queued still run on buffers that were never filled, and a payload can be torn. Clear the header before every recv, and put a checksum on the payload. Destroy the graph exec before ncclCommAbort.
The state that makes per-message exactly-once recovery work sits at the endpoint, not in the communication library. Sequence numbers, dedup, and a resend against the peer’s consumption log already cover the completion ambiguity seen here, at very low cost. What is still open is the boundary this protocol does not cover: whether consumption and the log have to be atomic, how to handle consumption that has an external side effect, and partial completion of a collective that cannot be written as one message.
Recovery should commit on what the receiver has used
The opening question was whether a given message should be resent when NCCL recovery runs. The experiment’s answer is not “read the error code,” and it is not “read sender completion.” It is to name the semantics you actually need to protect.
If the semantics are only “the send buffer may be reused now,” sender completion is enough. If they are “the bytes reached the peer,” you need delivery. If they are “this message has already changed receiver application state,” the commit point has to be consumption. Calling all three events completion is harmless on a healthy run. Under partial failure, the gap between them becomes a loss or a duplicate.
The result worth keeping is not the quirk of one NCCL API. It is a more general systems judgment. Failure recovery has to define commit around the application effect. It cannot treat transport or runtime completion as application completion. On these two V100s with EDR InfiniBand and no GPUDirect RDMA, that difference is already enough to produce a stable 1–4 missed resends or 4–9 duplicates. Once recovery queries the receiver’s consumption facts, both errors disappear.
The line that matters is who holds the last information about whether the operation really happened. In this experiment, that is the receiver application.
发送端看到 completion 之后,接收端可能还没有真正消费这条消息。故障恢复到底应该相信发送端的完成状态,还是接收端已经产生的应用效果?
这篇文章讨论的是一个很具体的问题:当 NCCL 通信进行到一半发生故障时,我们怎样判断一条消息该不该重发?
假设发送端发出一段 KV cache 或 expert output,并且已经看到 CUDA event 完成。此时接收端进程突然被 kill。恢复之后有两个选择:不重发,或者重发。如果不重发,前提是发送端 completion 已经足以证明接收端完成了处理;如果重发,前提是接收端还没有真正使用这条消息。问题在于,这两个前提可能同时都不可靠。
我们最初的猜想很简单:NCCL 暴露给发送端的 completion,加上 communicator 的错误状态,也许已经足够定义重试边界。 如果这个猜想成立,上层不需要维护额外的逐消息状态;恢复时只需要重发那些发送端没有看到 completion 的操作。
真正要验证的因此不是 NCCL 会不会出错,而是下面这条等价关系是否成立:
sender completion ≈ receiver has already produced the application effect
后面的实验基本都围绕这条关系展开。
为什么这个问题现在开始变重要
过去,大规模训练遇到通信故障时,常见处理方式是让整个通信组失败,然后从 checkpoint 重新开始。这样做虽然昂贵,但语义很简单:既然整步都丢掉,就不必追究某一条在途消息到底走到了哪里。
现在的系统越来越不愿意这样做。torchft 可以在 step 边界剔除失败副本,NCCL 提供 ncclCommShrink 与 ncclCommRevoke 让 communicator 在故障后继续重组。推理侧也越来越依赖长期状态:expert parallel 的 token dispatch、prefill/decode 分离中的 KV cache 传输,以及长时间 RL rollout,都希望保住已经完成的工作,而不是把整个 step 清空重来(Qin et al., FAST 2025)。
flowchart LR
A["过去: rank 失败"] --> B["整步作废"] --> C["从 checkpoint 重来"]
D["现在: rank 失败"] --> E["revoke / rebuild"] --> F["保留已有状态"]
F --> G["哪些消息已生效?"]
这张图是问题出现的根源。只要整个 step 都作废,消息级 completion 语义并不重要;一旦系统想保留部分进度,消息级语义立刻变成恢复正确性的基础。 例如一段 KV 已经被 decode worker 接入后续 token 生成,再重放一次可能造成重复状态;反过来,如果数据只是到了某个中间 buffer、还没被消费,却因为 sender completion 而不再补发,就会直接丢状态。
问题的核心:completion、delivery 和 consumption 不是一回事
读后文要先分清三个事件。第一层是本地完成:发送端已经把这次发送交给 NCCL/CUDA 通信流水线,发送缓冲现在可以安全复用。NCCL 的 ncclSend 本身是异步的,它只是把工作排到 CUDA stream;调用方通常在后面 record 一个 CUDA event,因此 event 完成代表发送侧这段工作走完了,而不是对方应用已经处理了。第二层是送达:字节真的落进接收端内存。第三层才是消费:接收端的 kernel 或应用逻辑读了这些字节,并把它们变成了不可忽略的效果,例如把一段 KV cache 接入后续 decode,或把 expert 输出累加进最终结果。
一个直观类比是快递状态:本地完成更像寄件柜已经收件,寄件人可以走了;送达是包裹到了收件地址;消费才是收件人已经拆包并使用里面的东西。正常运行时三者顺序稳定,很容易把它们统称为完成;故障恢复真正麻烦的地方就在于,故障可以恰好切在这三步之间。
flowchart LR
S["1. ncclSend"] --> C["2. Sender completion"] --> D["3. Delivery"] --> U["4. 应用效果"]
因此这篇文章真正关心的不是通信 API 名字,而是故障恢复的 commit point 到底在哪一层。如果 commit point 放在 sender completion,恢复可能太早相信一条消息已经完成;如果 completion 又是 graph 或 batch 粒度,恢复可能反过来太晚确认,导致已经消费的消息被重放。
在 RDMA 可靠连接(RC)上,一次 write 的完成事件意味着对端网卡已经确认数据写进对端内存,也就是送达,但仍不包括应用消费(InfiniBand 规范)。队列对(QP)进入错误状态后,未完成请求会以 flush error 返回,但其中一部分写操作可能早已在对端发生,所以请求失败也不能反推远端什么都没发生。
这套实验没有 GPUDirect RDMA,因此路径更长:网卡不能直接把远端数据写进 GPU HBM,而是先落到接收端 CPU 可见的 host bounce buffer,再由 NCCL proxy 和拷贝路径搬进 GPU,最后接收 kernel 与后续消费 kernel 才真正使用它。直觉上,这相当于在包裹到楼下和包裹送进房间并拆开之间又加了一层中转。每增加一个中转阶段,就多一个故障可能切进去的窗口;所以这里不是抽象意义上的语义争论,而是由真实数据路径直接产生的时间差。
现有的恢复原语都只给组粒度的信息。ncclCommAbort 丢弃所有在途操作;ncclCommRevoke 停掉在途操作并把 communicator 标成 in progress,静止后可以 destroy 或 shrink;ncclCommGetAsyncError 返回整个 communicator 的一个错误码(NCCL 文档)。PyTorch 的 ProcessGroupNCCL 在此之上加 watchdog,出错时让这个组上所有未完成的 Work 失败。MPI 的容错扩展 ULFM 提供 revoke、shrink 与 agree,agree 能让存活进程对一个整数达成容错共识(Bland et al., 2013),但出错的请求本身只返回进程失败或已撤销。UCX 的 peer 错误处理模式保证每个 send 请求都会以成功或错误结束,但错误结束不说明对端是否收到(UCX)。
量级上,这个窗口其实很大。本文环境里 abort 或 revoke 本身要约 0.5 s,主要来自 NCCL proxy 线程的 500 ms poll;而一条 16 B send/recv 往返只有几十微秒。也就是说,通信组从一端已经决定放弃到另一端真正静止的控制时间,比单条消息生命周期长了大约四个数量级。发送窗口虽然只允许最多 32 条操作同时在途,但这 32 条可以在这半秒里不断跨越已发送、已到达、已消费的边界。直觉上,问题不是网络突然变得很慢,而是恢复控制面太粗、太慢,而数据面仍在向前推进。
TL;DR
通信库告诉发送端我做完了,并不等于接收端应用已经用过这条消息。 在两台 V100、EDR InfiniBand、没有 GPUDirect RDMA 的 NCCL 2.32.3 上,我们做了 634 次故障注入。只按 sender completion 决定是否重发时,eager 和 revoke 每次故障会漏掉大约 1 到 4 条。这些消息发送端已经完成,接收端还没消费。16 条一批的 CUDA Graph 因为完成太粗,每次又会重复大约 4 到 9 条。把恢复依据换成接收端实际消费到的序号之后,重复和丢失都降到 0,额外开销落在测量噪声里。稳定的提交点是应用效果,不是传输层或运行时的 completion。
我们具体猜了什么
最简单的猜想是:sender completion 可以近似当作消息的恢复 commit point。 也就是把在途消息分成两类:发送端已经看到 completion 的,不再补发;还没看到 completion 的,故障恢复后重新发送。
这个猜想之所以看起来合理,是因为正常执行时 completion 和 consumption 往往紧挨着发生,而且 NCCL 已经替应用隐藏了 RDMA、proxy、buffer 和 stream 的很多细节。对上层来说,CUDA event 是最自然、也最容易获得的这次通信结束了的信号。
但它有两个可能的破口。第一,completion 可能太早:sender 已经认为发送结束,但 receiver 只是收到了数据,还没真正运行消费 kernel。第二,completion 可能太粗:CUDA Graph 把多个 send 合成一个 graph-level event,receiver 已经消费了一部分消息,sender 却仍然只能看到整个 graph 尚未完成。
所以我们同时放了一个更强的对照:每条消息带序号,接收端只有在真正消费后才把序号写进日志;恢复时发送端不看自己的 completion,而是查询接收端哪些序号真的已经产生应用效果,只补发缺失项。这对应经典的端到端原则:通信层能告诉你传输了多少,但只有端点知道应用语义上的操作是否真正发生(Saltzer, Reed, and Clark, 1984)。
实验真正比较的就是两个 commit point:
flowchart LR
A["故障"] --> B{"相信谁?"}
B --> C["Sender completion"]
B --> D["Receiver log"]
C --> C1["完成: 不补发"]
C --> C2["未完成: 补发"]
D --> D1["已消费: 不补发"]
D --> D2["未消费: 补发"]
C1 --> E["可能漏补"]
C2 --> F["batch 可能多补"]
D1 --> G["按应用效果"]
D2 --> G
如果 sender completion 与 receiver consumption 基本一致,两条路径应该得到相同结果;如果它们在故障窗口里系统性分叉,错误模式就会直接暴露出来。
实验如何把 completion、delivery 和 consumption 分开
图 1 是一条消息的生命周期。发送端 GPU 上一个 kernel 填好消息,ncclSend 发出,其后的 CUDA event 完成就是发送端看到的完成。数据经 IB RC 写进接收端 host 缓冲,再拷到 GPU,ncclRecv 完成后接收端的消费 kernel 校验并应用它,最后把序号追加进一个持久日志,我们以日志为准定义已消费。灰色区域是发送端已报告完成、接收端还没消费的窗口,故障落在这里时,只补发没完成的消息就会丢。图下方说的是 Graph 路径:16 次 send 被捕获进一个 CUDA graph,发送端只能看到整个 graph 的一个完成事件。
硬件是两台节点各一张 V100-PCIe-32GB,EDR InfiniBand,没有 GPUDirect RDMA,所以数据走 host bounce。软件是 NCCL 2.32.3、CUDA 12.4。每条消息的头部 16 B 是 64 位序号与校验和,负载每个字由序号决定,接收端逐条校验。消息大小在 16 B 与 64 KB 之间按 trial 交替,发送窗口 32 条,Graph 路径每个 graph 16 条、两个 graph 在途。一个独立的 TCP 控制连接传 communicator 的 uniqueId、故障计划、abort 通知与事后的日志查询。
三条恢复路径:eager 是逐条 send/recv,故障后 ncclCommAbort 再新建 communicator;Graph 同样用 abort;revoke 用 nonblocking communicator,故障后 ncclCommRevoke 加 destroy 再新建。三类故障:F1 在流中途 SIGKILL 接收进程,由 supervisor 重启并重新加入;F2 是发送端在流中途主动 abort;F3 是 SIGSTOP 接收进程 100 到 800 ms 再 SIGCONT,暂停连同 NCCL 的 proxy 线程一起。注入时刻在无故障流长度内均匀随机。检测只用现有信号:控制 socket 断开、ncclCommGetAsyncError、500 ms 无进度超时、以及 F2 的自我 abort。
我们比较两种恢复口径。第一种只使用发送端能观察到的 completion:补发所有发送端尚未看到完成的消息,接收端不去重。第二种使用端点消费状态:消息带序号,接收端用去重位图并在真正消费后写日志,进程重启后从日志恢复;通信恢复后,发送端查询接收端日志,只补发日志里缺失的序号。计数定义:重复是同一序号被消费两次以上;丢失是恢复结束后日志里仍没有的序号;stale 是故障时刻发送端完成状态与接收端消费状态不一致的条数。每格 210 次注入,F3 里短于 500 ms 的暂停不会触发任何检测,只增加延迟,所以触发次数少于 210。另跑无故障对照各 40 次,四个计数全 0,排除 harness 误报。第四类故障链路抖动需要改网卡端口状态,在共享节点上不允许,没有测。
发送端 completion 与接收端消费状态怎样分叉
图 2 的三个子图对应三类故障,横轴是恢复路径与策略,纵轴是每次真正触发恢复时平均出现多少条重复(实心)或丢失(斜线)消息。读这张图时不用先盯绝对值,更重要的是看错误方向:eager 与 revoke 主要向丢失偏,Graph 主要向重复偏。这两个方向恰好对应两种 completion 粒度:前者太早把单条发送判成完成,后者又太晚才把一整个 batch 判成完成。
eager 路径固定重试下,F1 接收端被杀时平均每次丢 4.0 条(169 次触发共 675 条),F3 暂停时约 1.0 条,F2 发送端 abort 时 171 次里丢 28 条、重复 1 条。相对于最多 32 条在途操作,4 条并不是整个窗口都错了,而是大约最后几条正好卡在发送端已经放行、接收端还没真正消费的缝里。这个数量级也符合数据路径直觉:绝大多数消息要么还没走到 completion,要么已经走到消费,只有流水线尾部的一小段处在模糊区。
这些丢失消息的共同点是发送端已经看到 completion,因此恢复逻辑不会重发;但接收端可能只拿到了 host 或 GPU buffer,还没运行消费 kernel,进程就被 kill,或者 abort 让这批尚未消费的数据失效。所以 eager 的错误不是重试太激进,而是发送 completion 作为确认点太早。 revoke 路径的形状相同:F1 平均每次丢 2.1 条,F2 与 F3 分别累计 27 与 91 条,没有重复。revoke 把 kill 路径从 eager 的约 880 ms 缩到 385–520 ms,改善的是多久能停下来,但没有改变哪一个事件算消息已经发生的语义。
Graph 路径正好反过来,以重复为主:F2 下 210 次触发共 1901 条重复,平均每次约 9 条;F1 与 F3 分别约 3.9 与 6.2 条。这里最直观的理解是:CUDA Graph 给发送端的是一张16 条消息的批量收据。故障发生时,接收端可能已经消费了前 9 条,但整张 graph 的完成事件还没亮,于是恢复逻辑只能把 16 条都当成没做完再发一次,前 9 条自然重复。
因此 9 条这个数字不是随机噪声,而是在 16-op batch 中途发生故障时很自然的量级:越靠近 batch 尾部发生故障,已经消费但尚未获得 batch completion 的消息越多。Graph 把 completion 变粗之后,错误从少发翻成了多发。 如果消费是非幂等的,例如累加梯度、合并 expert 输出或推进某个状态机,这类重复不是简单的性能浪费,而可能直接变成静默数值错误。
端到端协议那一栏在三类故障下全是 0。eager 路径 F1、F2、F3 各 178、177、68 次触发(F2 另补 4 次,共 214 次注入),重复与丢失都是 0。更有意思的是,它并没有用更复杂的通信原语:只是让接收端记录我真正消费到了哪些序号,恢复时按这份事实账本补缺口。发送端 completion 记录的是发送侧进度,接收端日志记录的是应用效果;一旦恢复依据从前者换成后者,前面两种相反的错误同时消失。
代价也很小。无故障流上,打开去重后 16 B 消息吞吐变化 +1.7%,64 KB 为 -0.8%,都在 run-to-run 噪声里。这里的直觉是,序号比较和位图查询只是很轻的 per-message 元数据操作,相比真正的数据搬运与 GPU kernel 调度并不占主导。Graph 路径我们只跑了固定重试;离线叠加接收端去重后,重复归零,剩下的 F1 121 条、F3 6 条缺口恰好都是接收端日志里没有的序号,也就是同一种按日志补发策略能够识别的部分,不过这一点没有在 Graph 在线恢复路径上单独实测。
另外三个直接暴露出来的 failure mode
第一,ncclCommGetAsyncError 在全部触发故障里都返回成功,包括对端进程已经被 SIGKILL。直觉上这说明peer 已经死了,与 communicator 已经形成一个可查询的异步错误,不是同一个时刻。实验里的故障检测实际上来自控制 socket 断开、500 ms 无进度超时,或发送端主动 abort。也就是说,如果上层把这个 API 当成类似 TCP connection reset 的即时死亡通知,在这个环境里就会一直等。
第二,abort 并不会把 GPU stream 里已经排好的后续 kernel 一并撤回。接收端仍有大量消费 kernel 被调度执行,只是它们对应的 recv slot 从未被新数据填充,每格能看到约 2000–6400 次这样的空读。最简单的直觉是:通信控制面停止了,但 GPU 执行队列不会自动倒带。 我们在每次 recv 前先把 16 B 头清零,因此这些 kernel 看到空头就拒绝消费;另外还有 14 条消息出现头部像新的、payload 却半新半旧的撕裂状态,被校验和抓到。没有显式的空值标记和完整性校验,这两类都会伪装成合法输入。
第三,资源销毁顺序本身会决定恢复能否结束:当 CUDA graph exec 仍然持有 communicator 时,ncclCommAbort 超过 35 s 不返回;先销毁 graph exec,再 abort 就恢复正常。直觉上这是一个生命周期依赖反转:你以为应该先停止通信再清理 graph,但 graph 本身仍持有通信资源,于是 abort 在等一个不会先释放的对象。对 Graph 化推理引擎来说,这比单纯的延迟问题更危险,因为错误的 cleanup 顺序会把可恢复故障变成恢复路径自身的 hang。
端点消费状态怎样消除这类歧义
图 2 最右一栏给出了一个很重要的边界:NCCL 暴露出来的 completion 与 communicator error 不足以决定逐消息恢复,但接收端的消费状态足以。在 634 次注入中,序号、去重和按接收端日志补发把重复与丢失都降到了 0;无故障流里,16 B 消息吞吐变化为 +1.7%,64 KB 为 -0.8%,都落在 run-to-run 噪声范围内。
这里真正起作用的信息不是更细的 NIC completion,也不是更快的 communicator teardown,而是应用已经消费到哪个序号。可以把两端维护的状态理解成两本账:发送端 completion 是执行流水账,我把哪些操作交出去了;接收端消费日志是效果账,哪些操作真的改变了我的应用状态。故障恢复真正需要的是后一本。只看前一本,就会在 eager 中漏补,在 Graph 中多补;按后一本补差额,两种错误同时消失。换句话说,partial failure 下的正确恢复边界天然跨过通信库,落到了端到端的应用效果状态上。
这个结果也限定了还没回答的问题。这里只有两个 rank 的点对点通信,集合通信的部分完成对不上逐条消息。环境没有 GPUDirect RDMA,host bounce 会把本地完成和消费之间的窗口拉大。去重检查和日志追加在同一个消费 kernel 里,我们没有在这两步之间再注入故障,所以带外部副作用的消费,原子性还要单独处理。链路抖动和 MSCCL++ 没有实测。MSCCL++ 0.9.0 没有 abort、revoke 和逐操作状态。源码里能看到它的恢复边界更粗,但这代替不了故障注入。
从这些数据可以直接得到什么
在 NCCL 2.32.3 的点对点 send/recv 上,发送端 CUDA event 完成只说明本地发送结束了。接收进程在故障里死掉,或者被 abort 时,每次故障有 1 到 4 条已经报告完成、却从未被消费的消息。这是在两张 V100、EDR IB、没有 GDR、窗口 32 条、16 B 与 64 KB 消息上测的。按发送端完成来决定补发什么,会把这些消息丢掉。
CUDA Graph 捕获的通信按整个 graph 报告完成。故障发生时,接收端可能已经消费了一部分。把整批还没看到完成的消息重发,每次故障会重复大约 4 到 9 条。这里一个 graph 是 16 条。粒度越粗,重复越多。Graph 化的通信要在每条消息上带序号,不能靠 graph 这一级的完成。
对端进程被杀、被暂停,或者本端主动 abort 时,ncclCommGetAsyncError 一直返回成功。故障检测要另找一条控制通道,或者在应用层做超时。
abort 之后,已经排进队列的接收 kernel 仍会在没填上数据的缓冲上跑,负载也可能被写到一半。每次 recv 之前要把消息头清掉,负载要带校验和。ncclCommAbort 之前必须先销毁 graph exec。
逐消息 exactly-once 恢复要看的状态在端点,不在通信库。序号、去重,再加上按对端消费日志补发,已经盖住了这次看到的 completion 歧义,开销也很小。后面更值得追的,是这套协议盖不住的边界:消费和写日志之间要不要做成原子,带外部副作用的消费怎么处理,以及集合通信里没法写成单条消息的部分完成。
恢复该按接收端已经用掉的内容提交
文章一开始问的是:NCCL 故障恢复时,一条消息到底该不该重发? 实验给出的答案不是看错误码,也不是看 sender completion,而是先明确你真正想保护的语义是什么。
如果语义只是发送缓冲现在能不能复用,sender completion 已经足够;如果语义是字节有没有到对端,需要看 delivery;如果语义是这条消息是否已经改变了接收端应用状态,那么 commit point 就必须落到 consumption。把这三个事件都叫 completion,正常运行时问题不大;一旦进入 partial failure,它们之间的缝就会直接变成丢失或重复。
这组数据最值得记住的不是某个 NCCL API 的具体行为,而是一个更一般的系统判断:故障恢复必须围绕 application effect 定义 commit,而不能把底层 transport 或 runtime 的 completion 直接等同于应用完成。 在我们的两机 V100 加 EDR IB、无 GDR 环境里,这个差异已经足以稳定地产生 1–4 条漏补或 4–9 条重复;而一旦恢复改为查询 receiver 的消费事实,两种错误同时消失。
所以这里真正的分界线不是通信库够不够可靠,而是谁拥有判断操作是否已经真正发生的最后信息。在这个实验里,答案是接收端应用。