Modal logo

Multi-node Clusters

Modal Clusters supports petaFLOP/s jobs that run across several coordinated containers, such as serving or training trillion-parameter models. Each container can saturate the available GPU devices on its node and communicate with peer containers for terabit/s networking.

Modal Clusters provide:

When you call a Clustered Function or Server, Modal starts multiple containers simultaneously, then runs your code in each container. Therefore, your code should establish communication between containers and coordinate with each other to deliver the final result.

The guide will walk you through how to fully take advantage of Modal Clusters.

Using the @clustered decorator 

To make a Clustered Function or Server, use the @clustered decorator:

@app.function(
    gpu="H100:8",
    timeout=60 * 60 * 24,
    retries=modal.Retries(initial_delay=0.0, max_retries=10),
)
@modal.clustered(size=4)
def train_model():
    cluster = modal.Cluster.from_context()

    container_rank = cluster.container_rank()
    world_size = len(cluster.container_ips())
    main_addr = cluster.container_ips()[0]
    is_main = "(main)" if container_rank == 0 else ""

    print(f"{container_rank=} {is_main} {world_size=} {main_addr=}")

The above configuration creates a group of four containers each having eight H100 GPU devices, for a total of 32 devices.

Containers in a multi-node cluster are physically colocated and gang scheduled together so that your code only runs once all of the requested hardware is acquired.

Traditionally this kind of cluster and scheduling management would be handled by SLURM, Kubernetes, or something else. But with Modal, it’s all provided serverlessly with just a Python decorator!

Rank & input broadcast 

Each container in a multi-node cluster is assigned a rank. Rank zero is the “leader” rank, or head node, and typically coordinates the job.

To determine the current container’s rank, use the modal.Cluster API:

cluster = modal.Cluster.from_context()
container_rank = cluster.container_rank()
if container_rank == 0:
    print("Running as the leader")
else:
    print(f"Running as rank {container_rank}")

When you call a Clustered Function, each container receives a copy of the Function call’s arguments. For example, if you call a four-node Function, your code will run four times in parallel across four containers.

For Clustered Servers, network traffic is routed only to rank zero, which is then responsible for distributing requests to the other containers.

Only rank zero’s output is returned to the caller; i.e., outputs from other ranks are discarded. To share work before returning a final result, use inter-container networking or RDMA.

Networking 

In addition to gang scheduling, the @clustered decorator enables i6pn, Modal’s workspace-private inter-container networking, so that containers in a cluster can communicate with each other over TCP or UDP. You can then use protocols such as Torch Distributed Elastic on top of Modal’s networking stack.

To get the list of IPs for each container in a cluster:

cluster = modal.Cluster.from_context()
container_ips = cluster.container_ips()
print(f"The main container's IP is {container_ips[0]}")

By default, container_ips will return IPv6 addresses. For workloads that require IPv4 (such as Ray-based frameworks), pass in the family parameter:

cluster = modal.Cluster.from_context()
container_ips = cluster.container_ips(family="ipv4")
print(f"The main container's IP is {container_ips[0]}")

For more on what you can do with inter-container networking, see the cluster networking guide.

RDMA 

For even higher inter-node bandwidth, you can enable RDMA. The exact bandwidth depends on GPU type: B300 clusters have 6,400 Gbps networking, while other GPU types have 3,200 Gbps.

To use RDMA, make sure your container image contains the necessary dependencies, typically a copy of libcudart.so, libibverbs.so.1, and libmlx5.so.1. The easiest way to do this is to use a CUDA base image, then .apt_install the InfiniBand library:

cuda_version = "12.9.1"
flavor = "devel"
operating_sys = "ubuntu22.04"
tag = f"{cuda_version}-{flavor}-{operating_sys}"

image = (
    modal.Image.from_registry(f"nvidia/cuda:{tag}", add_python="3.12")
    .apt_install("libibverbs1")
)

Then, pass in rdma=True to the clustered decorator:

@modal.clustered(size=2, rdma=True)
def train():
    ...

If you’re using NCCL or a framework based on it, Modal automatically sets the necessary environment variables in your containers.

Otherwise, Modal exposes one of two RDMA interfaces, depending on the underlying hardware: InfiniBand Verbs or EFA. If your workload is not compatible with EFA, you can force your container to run on an InfiniBand Verbs host using the efa_disabled experimental option:

@app.function(
    ...,
    experimental_options={
        "efa_disabled": True,
    },
)
@modal.clustered(size=2, rdma=True)
def train():
    ...

To run a simple RDMA performance test, see this sample code.

Fault tolerance 

Failures are propagated to other containers. If an input fails on any individual container, Modal will terminate all remaining containers and mark the entire call as failed, even if the input succeeded on another container. You can set a retry policy to have Modal retry the input for you.

Preemptions apply to the entire cluster. In the event of a preemption, Modal will terminate all containers in the cluster and retry it with the same input.

Input synchronization 

Modal does not synchronize input execution across containers. Containers are responsible for ensuring that they do not process inputs faster than other containers in their cluster.

In particular, it is important that the leader container (rank 0) waits for all other containers to finish processing the current input before starting the next one.