nemo_automodel.components.moe.uccl_ep.buffer

View as Markdown

UCCLBuffer: a DeepEP-compatible Buffer backed by UCCL-EP.

This module re-exports the canonical Buffer implementation under the UCCLBuffer alias expected by nemo_automodel, with automatic intranode detection.

Module Contents

Classes

NameDescription
BufferThe core expert-parallel (EP) communication buffers for Mixture of Experts (MoE) model, which supports:
EventHandle-
EventOverlapA wrapper class to manage CUDA events, also for better overlapping convenience.
UCCLBufferBuffer subclass that auto-detects intranode mode.

API

class nemo_automodel.components.moe.uccl_ep._buffer.Buffer(
group: torch.distributed.ProcessGroup,
num_nvl_bytes: int = 0,
num_rdma_bytes: int = 0,
low_latency_mode: bool = False,
num_qps_per_rank: int = 24,
allow_nvlink_for_low_latency_mode: bool = True,
allow_mnnvl: bool = False,
explicitly_destroy: bool = False,
is_intranode: bool = False
)

The core expert-parallel (EP) communication buffers for Mixture of Experts (MoE) model, which supports:

  • high-throughput intranode all-to-all (dispatch and combine, using NVLink)
  • high-throughput internode all-to-all (dispatch and combine, using RDMA and NVLink)
  • low-latency all-to-all (dispatch and combine, using RDMA)
group_size
= group.size()

the number of ranks in the group.

num_sms
int = 20

the SMs used in high-throughput kernels.

rank
= group.rank()

the local rank number.

runtime

the C++ runtime.

scratch
nemo_automodel.components.moe.uccl_ep._buffer.Buffer._dtype_code(
dtype: torch.dtype
) -> int
staticmethod
nemo_automodel.components.moe.uccl_ep._buffer.Buffer._ll_compute_stream_ptr(
device: torch.device
)

Return the current CUDA stream pointer for low-latency runtime calls.

nemo_automodel.components.moe.uccl_ep._buffer.Buffer._unpack_bias(
bias: typing.Union[torch.Tensor, typing.Tuple[torch.Tensor, torch.Tensor]]
)
staticmethod
nemo_automodel.components.moe.uccl_ep._buffer.Buffer.capture() -> nemo_automodel.components.moe.uccl_ep._utils.EventOverlap
staticmethod

Capture a CUDA event on the current stream, i.e. torch.cuda.current_stream().

Returns: EventOverlap

the captured event.

nemo_automodel.components.moe.uccl_ep._buffer.Buffer.clean_low_latency_buffer(
num_max_dispatch_tokens_per_rank: int,
hidden: int,
num_experts: int
) -> None

As low-latency kernels require part of the buffer to be zero-initialized, so it is vital to clean the buffer if the buffer is dirty at some time. For example, after running the normal dispatch/combine, you must run this function before executing any low-latency kernel.

Parameters:

num_max_dispatch_tokens_per_rank
int

the maximum number of tokens to dispatch, all the ranks must hold the same value.

hidden
int

the hidden dimension of each token.

num_experts
int

the number of all experts.

nemo_automodel.components.moe.uccl_ep._buffer.Buffer.combine(
x: torch.Tensor,
handle: typing.Tuple,
topk_weights: torch.Tensor | None = None,
bias: typing.Union[torch.Tensor, typing.Tuple[torch.Tensor, torch.Tensor]] = None,
config: ep.Config | None = None,
async_finish: bool = False,
allocate_on_comm_stream: bool = False
) -> typing.Tuple[torch.Tensor, torch.Tensor | None, nemo_automodel.components.moe.uccl_ep._utils.EventOverlap]

Combine (reduce) tokens (addition without weights) from different ranks, both intranode and internode settings are supported. Intranode kernels require all the ranks should be visible via NVLink. Internode kernels require the ranks in a node should be visible via NVLink, while the ranks with the same GPU index should be visible via RDMA.

Parameters:

x
torch.Tensor

[num_tokens, hidden] with torch.bfloat16, the tokens to send for reducing to its original ranks.

handle
Tuple

a must-set communication handle, you can obtain this from the dispatch function.

topk_weights
torch.Tensor | NoneDefaults to None

[num_tokens, num_topk] with torch.float, the tokens’ top-k weights for reducing to its original ranks.

config
Config | NoneDefaults to None

the performance tuning config.

previous_event
EventOverlap | NoneDefaults to None

the event to wait before actually executing the kernel.

async_finish
boolDefaults to False

the current stream will not wait for the communication kernels to be finished if set.

allocate_on_comm_stream
boolDefaults to False

control whether all the allocated tensors’ ownership to be on the communication stream.

Returns: torch.Tensor

the reduced token from its dispatched ranks.

nemo_automodel.components.moe.uccl_ep._buffer.Buffer.connect_atomic_buffer(
proxy: 'ep.UcclProxy'
)
nemo_automodel.components.moe.uccl_ep._buffer.Buffer.destroy()

Destroy the cpp runtime and release resources.

nemo_automodel.components.moe.uccl_ep._buffer.Buffer.dispatch(
x: typing.Union[torch.Tensor, typing.Tuple[torch.Tensor, torch.Tensor]],
handle: typing.Tuple | None = None,
num_tokens_per_rank: torch.Tensor | None = None,
num_tokens_per_rdma_rank: torch.Tensor | None = None,
is_token_in_rank: torch.Tensor | None = None,
num_tokens_per_expert: torch.Tensor | None = None,
topk_idx: torch.Tensor | None = None,
topk_weights: torch.Tensor | None = None,
expert_alignment: int = 1,
num_worst_tokens: int = 0,
config: ep.Config | None = None,
async_finish: bool = False,
allocate_on_comm_stream: bool = False
) -> typing.Tuple[typing.Union[typing.Tuple[torch.Tensor, torch.Tensor], torch.Tensor], torch.Tensor | None, torch.Tensor | None, typing.List[int], typing.Tuple, nemo_automodel.components.moe.uccl_ep._utils.EventOverlap]

Dispatch tokens to different ranks, both intranode and internode settings are supported. Intranode kernels require all the ranks should be visible via NVLink. Internode kernels require the ranks in a node should be visible via NVLink, while the ranks with the same GPU index should be visible via RDMA.

Parameters:

x
Union[torch.Tensor, Tuple[torch.Tensor, torch.Tensor]]

torch.Tensor or tuple of torch.Tensor, for the first type, the shape must be [num_tokens, hidden], and type must be torch.bfloat16; for the second type, the first element of the tuple must be shaped as [num_tokens, hidden] with type torch.float8_e4m3fn, the second must be [num_tokens, hidden // 128] (requiring divisible) with type torch.float.

handle
Tuple | NoneDefaults to None

an optional communication handle, if set, the CPU will reuse the layout information to save some time.

num_tokens_per_rank
torch.Tensor | NoneDefaults to None

[num_ranks] with torch.int, the number of tokens to be sent to each rank.

num_tokens_per_rdma_rank
torch.Tensor | NoneDefaults to None

[num_rdma_ranks] with torch.int, the number of tokens to be sent to each RDMA rank (with the same GPU index), return None for intranode settings.

is_token_in_rank
torch.Tensor | NoneDefaults to None

[num_tokens, num_ranks] with torch.bool, whether a token be sent to a rank.

num_tokens_per_expert
torch.Tensor | NoneDefaults to None

[num_experts] with torch.int, the number of tokens to be sent to each expert.

topk_idx
torch.Tensor | NoneDefaults to None

[num_tokens, num_topk] with torch.int64, the expert indices selected by each token, -1 means no selections.

topk_weights
torch.Tensor | NoneDefaults to None

[num_tokens, num_topk] with torch.float, the expert weights of each token to dispatch.

expert_alignment
intDefaults to 1

align the number of tokens received by each local expert to this variable.

num_worst_tokens
intDefaults to 0

the worst number of tokens to receive, if specified, there will be no CPU sync, and it will be CUDA-graph compatible. Please also notice that this flag is for intranode only.

config
Config | NoneDefaults to None

the performance tuning config.

previous_event
EventOverlap | NoneDefaults to None

the event to wait before actually executing the kernel.

async_finish
boolDefaults to False

the current stream will not wait for the communication kernels to be finished if set.

allocate_on_comm_stream
boolDefaults to False

control whether all the allocated tensors’ ownership to be on the communication stream.

Returns: Union[Tuple[torch.Tensor, torch.Tensor], torch.Tensor]

received tokens, the same type and tuple as the input x, but the number of tokens equals to the received token count.

nemo_automodel.components.moe.uccl_ep._buffer.Buffer.get_combine_config(
num_ranks: int
) -> ep.Config
staticmethod

Get a recommended combine config.

Argument

num_ranks: the number of ranks.

Returns: Config

the recommended config.

nemo_automodel.components.moe.uccl_ep._buffer.Buffer.get_comm_stream() -> torch.Stream

Get the communication stream.

Returns: torch.Stream

the communication stream.

nemo_automodel.components.moe.uccl_ep._buffer.Buffer.get_dispatch_config(
num_ranks: int
) -> ep.Config
staticmethod

Get a recommended dispatch config.

Argument

num_ranks: the number of ranks.

Returns: Config

the recommended config.

nemo_automodel.components.moe.uccl_ep._buffer.Buffer.get_dispatch_layout(
topk_idx: torch.Tensor,
num_experts: int,
async_finish: bool = False,
allocate_on_comm_stream: bool = False
) -> typing.Tuple[torch.Tensor, torch.Tensor | None, torch.Tensor, torch.Tensor, nemo_automodel.components.moe.uccl_ep._utils.EventOverlap]

Calculate the layout required for later communication.

Parameters:

topk_idx
torch.Tensor

[num_tokens, num_topk], dtype must be torch.int64, the expert indices selected by each token, -1 means no selections.

num_experts
int

the number of experts.

previous_event
EventOverlap | NoneDefaults to None

the event to wait before actually executing the kernel.

async_finish
boolDefaults to False

the current stream will not wait for the communication kernels to be finished if set.

allocate_on_comm_stream
boolDefaults to False

control whether all the allocated tensors’ ownership to be on the communication stream.

Returns: torch.Tensor

[num_ranks] with torch.int, the number of tokens to be sent to each rank.

nemo_automodel.components.moe.uccl_ep._buffer.Buffer.get_local_buffer_tensor(
dtype: torch.dtype,
size: torch.Size | None = None,
offset: int = 0,
use_rdma_buffer: bool = False
) -> torch.Tensor

Get the raw buffer (slice supported) as a PyTorch tensor.

Argument: dtype: the data type (PyTorch dtype) for the tensor. size: the slice size (by elements) to get from the buffer. offset: the offset of the beginning element. use_rdma_buffer: whether to return the RDMA buffer.

nemo_automodel.components.moe.uccl_ep._buffer.Buffer.get_low_latency_rdma_size_hint(
num_max_dispatch_tokens_per_rank: int,
hidden: int,
num_ranks: int,
num_experts: int
) -> int
staticmethod

Get a minimum size requirement for the RDMA buffer. The size calculation will be done with BF16.

Parameters:

num_max_dispatch_tokens_per_rank
int

the maximum number of tokens to dispatch, all the ranks must hold the same value.

hidden
int

the hidden dimension of each token.

num_ranks
int

the number of EP group ranks.

num_experts
int

the number of all experts.

Returns: int

the RDMA buffer size recommended.

nemo_automodel.components.moe.uccl_ep._buffer.Buffer.get_next_low_latency_combine_buffer(
handle: object
)

Get the raw registered RDMA buffer tensor for next low-latency combine, so that the next combine kernel can skip the copying.

Parameters:

handle
object

the communication handle given by the dispatch function.

Returns:

the raw RDMA low-latency buffer as a BF16 PyTorch tensor with shape [num_local_experts, num_ranks * num_max_dispatch_tokens_per_rank, hidden], you should fill this buffer by yourself.

nemo_automodel.components.moe.uccl_ep._buffer.Buffer.internode_combine(
x: torch.Tensor,
handle: typing.Union[tuple, list],
topk_weights: torch.Tensor | None = None,
bias: typing.Union[torch.Tensor, typing.Tuple[torch.Tensor, torch.Tensor]] = None,
config: ep.Config | None = None,
async_finish: bool = False,
allocate_on_comm_stream: bool = False
) -> typing.Tuple[torch.Tensor, torch.Tensor | None, nemo_automodel.components.moe.uccl_ep._utils.EventOverlap]

Internode combine implementation, for more details, please refer to the combine docs. Normally, you should not directly call this function.

nemo_automodel.components.moe.uccl_ep._buffer.Buffer.internode_dispatch(
x: typing.Union[torch.Tensor, typing.Tuple[torch.Tensor, torch.Tensor]],
handle: typing.Tuple | None = None,
num_tokens_per_rank: torch.Tensor | None = None,
num_tokens_per_rdma_rank: torch.Tensor | None = None,
is_token_in_rank: torch.Tensor | None = None,
num_tokens_per_expert: torch.Tensor | None = None,
topk_idx: torch.Tensor | None = None,
topk_weights: torch.Tensor | None = None,
expert_alignment: int = 1,
num_worst_tokens: int = 0,
config: ep.Config | None = None,
async_finish: bool = False,
allocate_on_comm_stream: bool = False
) -> typing.Tuple[typing.Union[typing.Tuple[torch.Tensor, torch.Tensor], torch.Tensor], torch.Tensor | None, torch.Tensor | None, typing.List[int], typing.Tuple, nemo_automodel.components.moe.uccl_ep._utils.EventOverlap]

Internode dispatch implementation, for more details, please refer to the dispatch docs. Normally, you should not directly call this function.

nemo_automodel.components.moe.uccl_ep._buffer.Buffer.is_sm90_compiled()
staticmethod
nemo_automodel.components.moe.uccl_ep._buffer.Buffer.low_latency_combine(
x: torch.Tensor,
topk_idx: torch.Tensor,
topk_weights: torch.Tensor,
handle: tuple,
use_logfmt: bool = False,
zero_copy: bool = False,
async_finish: bool = False,
return_recv_hook: bool = False,
out: torch.Tensor | None = None,
combine_wait_recv_cost_stats: torch.Tensor | None = None
) -> typing.Tuple[torch.Tensor, nemo_automodel.components.moe.uccl_ep._utils.EventOverlap, typing.Callable]

A low-latency implementation for combining tokens (reduce with weights) with IBGDA. This kernel requires all the ranks (no matter intranode or internode) should be visible via RDMA (specifically, IBGDA must be enabled). Warning: as there are only two buffers, and the returned tensors reuse the buffer, you cannot hold more than 2 low-latency kernels’ result tensors at a single moment.

Parameters:

x
torch.Tensor

[num_local_experts, num_max_dispatch_tokens_per_rank * num_ranks, hidden] with torch.bfloat16, the local calculated tokens to be sent to this original rank and reduced.

topk_idx
torch.Tensor

[num_combined_tokens, num_topk] with torch.int64, the expert indices selected by the dispatched tokens. -1 indices (not selecting any expert) are supported. Note that, num_combined_tokens equals to the number of dispatched tokens.

topk_weights
torch.Tensor

[num_combined_tokens, num_topk] with torch.float, the expert weights selected by the dispatched tokens. The received tokens will be reduced with the weights in this tensor.

handle
tuple

the communication handle given by the dispatch function.

use_logfmt
boolDefaults to False

whether to use an internal “LogFMT with dynamic per-64-channel cast” format (10 bits).

zero_copy
boolDefaults to False

whether the tensor is already copied into the RDMA buffer, should be cooperative with get_next_low_latency_combine_buffer.

async_finish
boolDefaults to False

the current stream will not wait for the communication kernels to be finished if set.

return_recv_hook
boolDefaults to False

return a receiving hook if set. If set, the kernel will just do the RDMA request issues, but without actually receiving the data. You must call the received hook to make sure the data’s arrival. If you do not set this flag, the kernel will ensure the data’s arrival.

out
torch.Tensor | NoneDefaults to None

the in-place output tensor, if set, the kernel will write the result to this tensor and return it directly.

combine_wait_recv_cost_stats
torch.Tensor | NoneDefaults to None

a cumulative time spent waiting to receive each token tensor for statistics, which should have shape [num_ranks, num_ranks] and be typed as torch.int64. This is useful for detecting and pre-cisely localizing slow anomalies.

Returns: torch.Tensor

the reduced token tensor, with shape [num_combined_tokens, hidden] and type torch.bfloat16.

nemo_automodel.components.moe.uccl_ep._buffer.Buffer.low_latency_dispatch(
x: torch.Tensor,
topk_idx: torch.Tensor,
num_max_dispatch_tokens_per_rank: int,
num_experts: int,
cumulative_local_expert_recv_stats: torch.Tensor | None = None,
dispatch_wait_recv_cost_stats: torch.Tensor | None = None,
use_fp8: bool = True,
round_scale: bool = False,
use_ue8m0: bool = False,
async_finish: bool = False,
return_recv_hook: bool = False
) -> typing.Tuple[typing.Tuple[torch.Tensor, torch.Tensor], torch.Tensor, typing.Tuple, nemo_automodel.components.moe.uccl_ep._utils.EventOverlap, typing.Callable]

A low-latency implementation for dispatching with IBGDA. This kernel requires all the ranks (no matter intranode or internode) should be visible via RDMA (specifically, IBGDA must be enabled). Warning: as there are only two buffers, and the returned tensors reuse the buffer, you cannot hold more than 2 low-latency kernels’ result tensors at a single moment.

Parameters:

x
torch.Tensor

torch.Tensor with torch.bfloat16, shaped as [num_tokens, hidden], only several hidden shapes are supported. The number of tokens to be dispatched must be less than num_max_dispatch_tokens_per_rank.

topk_idx
torch.Tensor

torch.Tensor with torch.int64, shaped as [num_tokens, num_topk], only several top-k shapes are supported. -1 indices (not selecting any expert) are supported.

num_max_dispatch_tokens_per_rank
int

the maximum number of tokens to dispatch, all the ranks must hold the same value.

num_experts
int

the number of all experts.

cumulative_local_expert_recv_stats
torch.Tensor | NoneDefaults to None

a cumulative expert count tensor for statistics, which should have shape [num_local_experts] and be typed as torch.int. This is useful for online service EP load balance monitoring.

dispatch_wait_recv_cost_stats
torch.Tensor | NoneDefaults to None

a cumulative time spent waiting to receive each token tensor for statistics, which should have shape [num_ranks, num_ranks] and be typed as torch.int64. This is useful for detecting and pre-cisely localizing slow anomalies.

use_fp8
boolDefaults to True

whether to enable FP8 casting, with this, the received data will be a tuple of FP8 tensor and scaling factors.

round_scale
boolDefaults to False

whether round the scaling factors into power of 2.

use_ue8m0
boolDefaults to False

whether use UE8M0 as scaling factor format (available only with round_scale=True).

async_finish
boolDefaults to False

the current stream will not wait for the communication kernels to be finished if set.

return_recv_hook
boolDefaults to False

return a receiving hook if set. If set, the kernel will just do the RDMA request issues, but without actually receiving the data. You must call the received hook to make sure the data’s arrival. If you do not set this flag, the kernel will ensure the data’s arrival.

Returns: Tuple[torch.Tensor, torch.Tensor]

a tensor or tuple with received tokens for each expert. With use_fp8=True: the first element is a torch.Tensor shaped as [num_local_experts, num_max_dispatch_tokens_per_rank * num_ranks, hidden] with torch.float8_e4m3fn. The second tensor is the corresponding scales for the first element with shape [num_local_experts, num_max_dispatch_tokens_per_rank * num_ranks, hidden // 128] with torch.float, if use_ue8m0=False. With use_ue8m0=True, the second one is packed and shaped as [num_local_experts, num_max_dispatch_tokens_per_rank * num_ranks, hidden // 512] with type torch.int. Notice that, the last-two-dimension of the scaling tensors are in column-major for TMA compatibility. With use_fp8=False, the result would be a tensor shaped as [num_local_experts, num_max_dispatch_tokens_per_rank * num_ranks, hidden] with torch.bfloat16. Moreover, not all tokens are valid, only some of the num_max_dispatch_tokens_per_rank * num_ranks are, as we do not synchronize CPU received count with GPU (also not incompatible with CUDA graph if synced).

nemo_automodel.components.moe.uccl_ep._buffer.Buffer.reset_rdma_buffer()

Reset the RDMA buffer, this is useful when you want to reuse the RDMA buffer for another run.

nemo_automodel.components.moe.uccl_ep._buffer.Buffer.set_num_sms(
new_num_sms: int
) -> None
staticmethod

Set the number of SMs to use in high-throughput kernels.

Parameters:

new_num_sms
int

the new number to be set.

class nemo_automodel.components.moe.uccl_ep.buffer.EventHandle()
class nemo_automodel.components.moe.uccl_ep._utils.EventOverlap(
event: ep.EventHandle | None = None,
extra_tensors: typing.Tuple[torch.Tensor] | None = None
)

A wrapper class to manage CUDA events, also for better overlapping convenience.

nemo_automodel.components.moe.uccl_ep._utils.EventOverlap.__enter__() -> typing.Any

Utility for overlapping and Python with syntax.

You can overlap the kernels on the current stream with the following example:

event_overlap = event_after_all_to_all_kernels()
with event_overlap():
do_something_on_current_stream()
# After exiting the `with` scope, the current stream with wait the event to be finished.
nemo_automodel.components.moe.uccl_ep._utils.EventOverlap.__exit__(
exc_type: typing.Any,
exc_val: typing.Any,
exc_tb: typing.Any
) -> None

Utility for overlapping and Python with syntax.

Please follow the example in the __enter__ function.

nemo_automodel.components.moe.uccl_ep._utils.EventOverlap.current_stream_wait() -> None

The current stream torch.cuda.current_stream() waits for the event to be finished.

class nemo_automodel.components.moe.uccl_ep.buffer.UCCLBuffer(
group,
num_nvl_bytes: int = 0,
num_rdma_bytes: int = 0,
low_latency_mode: bool = False,
num_qps_per_rank: int = 24,
allow_nvlink_for_low_latency_mode: bool = True,
allow_mnnvl: bool = False,
explicitly_destroy: bool = False,
is_intranode: bool = False
)

Bases: Buffer

Buffer subclass that auto-detects intranode mode.

When all EP ranks fit on a single node (group_size <= LOCAL_WORLD_SIZE), RDMA is disabled and only NVLink is used, avoiding RDMA MR registration failures on single-node setups.