MSCCL, the Microsoft Collective Communication Library, is a runtime that executes custom collective algorithms on GPUs. It is built on top of NCCL: it keeps NCCL's API, transports and protocols, and adds an interpreter that runs a schedule you supply as an XML file. Its companion toolkit, msccl-tools, provides MSCCLang, a Python-embedded language for writing those schedules, and a compiler that turns them into the XML. When no custom algorithm matches a call, MSCCL falls back to NCCL's own ring and tree algorithms.

NCCL's built-in algorithms are general and near bandwidth-optimal for large messages, but not always optimal for a particular topology and size. A ring all-reduce over 8 GPUs takes 14 sequential steps, and on small messages those step latencies dominate. A full NVSwitch lets every GPU reach every other GPU directly, which allows a two-step schedule. MSCCL lets you write that schedule and register it for the size range where it wins.

Below: the architecture, a real MSCCLang program, the XML format, the rule that decides whether your algorithm runs, and a worked crossover example. For the primitives, start with collective operations and NCCL collectives.

Why generic collectives leave time on the table

In the alpha-beta model each synchronised step costs a fixed latency alpha (synchronisation, flag polling, link latency) and each byte costs 1/B, where B is the usable bandwidth. A ring all-reduce over p ranks is a reduce-scatter followed by an all-gather, each taking p-1 steps, and every rank sends (p-1)/p of the buffer in each phase:

T_ring(n)     = 2(p-1) * alpha  +  2 (p-1)/p * n / B
T_allpairs(n) = 2      * alpha  +  2 (p-1)/p * n / B     # needs a direct link to every peer

The bandwidth terms are identical. Both schedules move the minimum data a reduce-scatter plus all-gather needs. The latency terms are not. With p = 8, the ring pays 14 alphas and the all-pairs schedule pays 2. All-pairs is only possible when every rank can send to every other rank at full rate at the same time. A single node with NVSwitch provides that. A PCIe tree or a multi-node fabric with one NIC per GPU does not, and there the schedule turns into congestion.

That is the core MSCCL idea: the best algorithm depends on topology, size and protocol. The project's own example, an all-pairs all-reduce on an 8x A100 node, is described in the MSCCL README as being up to 2 to 3 times faster than the default ring for small and medium buffers. Treat that as the authors' claim for that machine. Measure your own.

Architecture: language, compiler, IR, runtime

MSCCL has four pieces. The MSCCLang program is Python that describes how chunks of each rank's buffer move and combine. The compiler lowers that program into per-GPU, per-thread-block instruction lists, and its Check() call verifies that the program actually computes the declared collective. The XML IR is the compiled schedule. The runtime is a fork of NCCL that loads the XML at communicator creation and runs it through an interpreter kernel. A second source of XML is synthesis. The project began as SCCL (PPoPP 2021, which searched for optimal algorithms with a solver), and TACCL (NSDI 2023) guides synthesis with user-written communication sketches. The compiler is described in the GC3 paper (ASPLOS 2023).

MSCCL: from a Python program to a schedule interpreted inside an NCCL-compatible runtimeMSCCLang programchunk(), copy(), reduce()Compiler + Check()instances, protocol, tb policyXML IR filealgo / gpu / tb / stepSynthesizer (SCCL, TACCL)solver output, same IRalternative sourceMSCCL runtime (NCCL fork)MSCCL_XML_FILES / MSCCL_CONFIG loaded at communicator initparsed per rankFramework callncclAllReduce(...)Selection predicateop, coll, inplace, ranks, sizeInterpreter kernelthread blocks run stepsmatchno match: NCCL's own ring / tree path runs, silently
The MSCCL pipeline. Programs and synthesizers both emit the same XML IR. The runtime checks a selection predicate on every collective call, and when it fails NCCL's built-in path runs with no error.

The runtime reuses NCCL's channels, transports and protocols (Simple, LL, LL128). It replaces only the decision about who sends which chunk to whom, in what order, on which thread block.

Writing an algorithm in MSCCLang

An MSCCLang program declares a topology, a collective (which fixes the pre- and postconditions that Check() verifies) and a number of instances. It then manipulates chunk(rank, buffer, index, size) references. c.copy(dst_rank, buffer, index) moves a chunk. c.reduce(other_chunk) combines one chunk into another. This is the all-pairs all-reduce from the msccl-tools examples, trimmed to the logic:

from msccl.language import *
from msccl.topologies import fully_connected
from msccl.language.collectives import AllReduce

def allreduce_allpairs(size, instances, protocol):
    collective = AllReduce(size, size * size, True)        # chunks per loop, in-place
    with MSCCLProgram("allreduce_pairs", fully_connected(size), collective,
                      instances, protocol=protocol,
                      threadblock_policy=ThreadblockPolicy.manual):
        # 1. scatter: rank r1 sends its r2-th slice straight to rank r2's scratch
        for r1 in range(size):
            for r2 in range(size):
                if r1 != r2:
                    c = chunk(r1, Buffer.input, r2 * size, size=size)
                    c.copy(r2, 'scratch', sendtb=r2, recvtb=r1)
        # 2. local reduction of everything that landed in scratch
        for r in range(size):
            for i in range(size * (size - 1)):
                c = chunk(r, Buffer.input, r * size + (i % size))
                c.reduce(chunk(r, 'scratch', i), sendtb=(i % size))
        # 3. gather: each rank broadcasts its fully reduced slice to every peer
        for r1 in range(size):
            for r2 in range(size):
                if r1 != r2:
                    c = chunk(r1, Buffer.input, r1 * size, size)
                    c.copy(r2, Buffer.input, r1 * size, sendtb=r2, recvtb=r1)
        XML()      # emit the IR on stdout
        Check()    # prove the result equals an all-reduce

Three details matter. First, size * size chunks per loop means the buffer is cut into 64 pieces on 8 GPUs. That becomes a divisibility rule at run time. Second, the manual thread-block policy gives each peer its own thread block via sendtb and recvtb. That is how the schedule drives all seven NVLink paths at once. Third, instances replicates the schedule across channels, buying parallelism with thread blocks taken from compute. Compile with python allreduce_a100_allpairs.py --protocol=LL 8 2 > allpairs.xml.

The XML IR

The runtime's parser defines the format. A root <algo> element carries name, ngpus, nchunksperloop, nchannels, proto (Simple, LL or LL128), coll (allreduce, allgather, reduce, broadcast, alltoall, reduce_scatter or custom), inplace and an optional minBytes/maxBytes size window. Inside it, each <gpu> declares its input, output and scratch chunk counts. Each <tb> (thread block) names one send peer, one receive peer and a channel. Each <step> is one instruction:

<!-- abridged and illustrative: attribute names from the parser, values hand-written -->
<algo name="allreduce_pairs" proto="LL" nchannels="2" nchunksperloop="64"
      ngpus="8" coll="allreduce" inplace="1" minBytes="0" maxBytes="4194304">
  <gpu id="0" i_chunks="64" o_chunks="0" s_chunks="56">
    <tb id="0" send="1" recv="1" chan="0">
      <step s="0" type="s" srcbuf="i" srcoff="8" dstbuf="s" dstoff="0"
            cnt="8" depid="-1" deps="-1" hasdep="0"/>
      ...
    </tb>
  </gpu>
  ...
</algo>

Step types are s (send), r (receive), rcs (receive, copy, send), rrs (receive, reduce, send), rrc (receive, reduce, copy), rrcs, cpy (local copy) and re (local reduce). depid/deps name a step on another thread block that must finish first. Hard limits in msccl.h include 256 steps per thread block and 32 thread blocks per channel. Scratch is sized as maxBytes * s_chunks / nchunksperloop, so an unbounded maxBytes costs real HBM.

How the runtime decides to use your algorithm

This is the most important section for anyone deploying MSCCL. On every collective call, the runtime's tuning code checks a predicate, and your algorithm runs only if all of these hold:

  1. The reduction op is sum, prod, max or min.
  2. The algorithm's coll matches the call (an all-gather XML never serves an all-reduce).
  3. The call's in-place-ness equals the XML's inplace. In-place means send buffer and receive buffer are the same pointer.
  4. ngpus equals the communicator's rank count. A sub-communicator of 4 ranks never matches an 8-GPU algorithm.
  5. The element count is divisible by nchunksperloop.
  6. minBytes <= nBytes < maxBytes.

If any check fails, NCCL's ring or tree runs and nothing is reported. The log line Connected 1 MSCCL algorithms proves only that the file parsed at init. It does not prove that any call used it. Algorithms are loaded with MSCCL_XML_FILES or MSCCL_CONFIG (an XML of <load> entries), and NCCL's algorithm list must allow it, for example NCCL_ALGO=MSCCL,RING,TREE. The README's benchmark compares nccl-tests' in-place column (the in-place MSCCL algorithm) against the out-of-place column (the default ring). That is predicate rule 3 in action. Confirm your framework's in-place behaviour with NCCL_DEBUG=INFO and a trace rather than assuming it.

Worked example: all-pairs versus ring on 8 GPUs

Plug numbers into the model for an 8-GPU NVSwitch node. Use about 300 GB/s per direction per GPU and an assumed 2 microseconds per synchronised step. Then account for protocol efficiency. LL writes 8-byte words that carry 4 data bytes and a 4-byte flag, so it gets about half the link rate in exchange for the lowest latency. LL128 carries 120 data bytes per 128-byte line, about 94%. The all-pairs schedule uses LL as in the example. The ring uses LL128.

MessageRing, LL128All-pairs, LLRing / all-pairs
64 KB28.4 us4.8 us5.96x
256 KB29.6 us7.1 us4.20x
1 MB34.5 us16.2 us2.13x
4 MB54.1 us52.9 us1.02x
16 MB132.4 us199.7 us0.66x
64 MB445.6 us786.9 us0.57x
Modelled all-reduce time on 8 GPUs (illustrative alpha-beta model, not a measurement)ring, LL128all-pairs, LL64 KB256 KB1 MB4 MB16 MB64 MBcrossover ~4.2 MBtime (log)message size (log scale)
The latency saving (12 alphas) is fixed. The bandwidth penalty of LL grows with size, so the curves cross.

Under these assumptions all-pairs wins by about 5.96x at 64 KB. The lines cross near 4.2 MB, and at 64 MB the ring is about 1.8 times faster. That is why MSCCL algorithms carry size windows. You register the custom schedule with maxBytes near the crossover and let NCCL take larger messages. The model is a planning tool, not a result: real alpha varies, and seven-peer contention is not free. Set the window from a measured nccl-tests sweep.

Small collectives show up in tensor-parallel all-reduces at small micro-batches, MoE token exchanges and decode-time all-reduces in tensor-parallel inference. Large data-parallel gradient buckets are bandwidth-bound, and a custom schedule rarely helps them.

Deploying it: nccl-tests, PyTorch and RCCL

To benchmark, build the runtime (make -j src.build), build nccl-tests against it with NCCL_HOME, and run all_reduce_perf with LD_LIBRARY_PATH pointing at MSCCL's build/lib. To use it in PyTorch, the framework has to load MSCCL's libnccl instead of its own. The msccl-tools README does this by replacing PyTorch's NCCL submodule and rebuilding, and its recipe pins PyTorch v1.9.0. Check with ldd that your build links NCCL dynamically, or the library path swaps nothing. msccl-tools has seen little activity since 2024, so check which NCCL version the fork tracks.

On AMD, the route is easier. RCCL integrates MSCCL, and its documentation says MSCCL is enabled by default on MI300X. On other platforms it needs RCCL_MSCCL_FORCE_ENABLE=1, and by default it is used only when each rank is its own process (RCCL_MSCCL_ENABLE_SINGLE_PROCESS=1 lifts that). MSCCL++ is a separate, newer project: a set of lower-level GPU communication primitives, not an XML interpreter. RCCL exposes it through RCCL_MSCCLPP_ENABLE.

# Prove the custom algorithm actually ran: compare runs with and without it
mpirun -np 8 -x LD_LIBRARY_PATH=msccl/build/lib:$LD_LIBRARY_PATH \
  -x NCCL_DEBUG=INFO -x NCCL_DEBUG_SUBSYS=INIT,ENV \
  -x MSCCL_XML_FILES=allpairs.xml -x NCCL_ALGO=MSCCL,RING,TREE \
  nccl-tests/build/all_reduce_perf -b 64K -e 64M -f 2 -g 1 -c 1 -n 100 -w 20

# Control run: same binary, MSCCL removed from the algorithm list
mpirun -np 8 ... -x NCCL_ALGO=RING,TREE nccl-tests/build/all_reduce_perf ...

Failure modes

The failure modes, roughly in order of how often they bite:

  • Silent fallback. One predicate rule fails, such as an odd tensor size, an out-of-place call or a 4-rank group. You measure NCCL and conclude MSCCL does not help. Always run an A/B with MSCCL removed from NCCL_ALGO. Identical curves mean it never ran.
  • Divisibility. With 64 chunks per loop, a count that is not a multiple of 64 elements falls back. Pad gradient buckets, or choose an algorithm with fewer chunks per loop.
  • Deadlock or wrong results from hand-edited XML. A missing dependency or a mismatched send/receive pair hangs the kernel. Only ship XML produced by a compiler run that passed Check(). Never hand-edit it.
  • SM theft. More instances and channels mean more thread blocks taking SMs from overlapped compute. A faster collective can make the step slower. Measure step time, not collective bus bandwidth. See overlapping collectives with compute.
  • Topology drift. A schedule for a fully connected node runs badly on one with a degraded link. Tie each XML to a hardware SKU.
  • Numerics. A different reduction order is not bitwise equal to the ring, which breaks bit-exact regression tests.

When MSCCL is worth it

MSCCL is a scalpel. Use it when you have a fixed fleet SKU, a profile showing collectives in a latency-bound size band on the critical path, and engineers who can own the XML across driver and library upgrades. Before writing one, check whether stock NCCL already closed the gap. Recent NCCL releases add an NVLS algorithm using NVLink SHARP in-switch reduction on supporting NVSwitch systems, so compare against the newest NCCL. Across nodes the bottleneck is usually the InfiniBand fabric, where hierarchical or TACCL-style schedules matter more. Inside a node, read NVLink Switch to understand which all-to-all patterns the fabric really supports.

What to do next

  1. Profile one training or inference step and list every collective with its size, rank count and whether it is in place.
  2. Mark the ones on the critical path between 64 KB and a few MB. If there are none, stop here.
  3. Baseline that size band with nccl-tests on the newest stock NCCL.
  4. Compile the all-pairs or ring example from msccl-tools for your GPU count and run Check().
  5. Benchmark with and without MSCCL in NCCL_ALGO and find the measured crossover.
  6. Set minBytes/maxBytes from that crossover and pad buffers to a multiple of nchunksperloop.
  7. Confirm on a real job that the algorithm is selected, and that step time, not only bus bandwidth, improves.
  8. Pin the XML to the hardware SKU and library version, and re-run the A/B after every NCCL or driver upgrade.
Key takeaway: MSCCL lets you replace NCCL's generic ring and tree with a schedule written for your topology and message size, compiled from MSCCLang into an XML IR and interpreted inside an NCCL-compatible runtime. The win is mostly latency on small and medium messages, it is bounded by a size window, and it only happens when every selection rule matches, so always prove the algorithm ran with an A/B before trusting a number.