Ray is a Python framework for running functions and classes across a cluster as if they were local. On GPU machines it has become a common glue layer: the scheduler that decides which process gets which GPU, the pipeline that streams data from storage into those GPUs, and the supervisor that restarts a training job when a node disappears. Libraries such as Ray Data, Ray Train and Ray Serve sit on top of one small core API, and nearly every production problem with them traces back to how that core treats GPUs.

This article explains that core from first principles: tasks, actors and objects; what num_gpus actually does (and does not do); placement groups for gang scheduling; why GPU tensors take a detour through host memory; and how batch inference and data-parallel training are built on these pieces. It ends with failure modes, trade-offs against plain torchrun and Kubernetes, and a checklist. Library APIs here were checked against the current Ray documentation; Ray moves quickly, so pin your version and re-read the docs for that version before copying flags.

Tasks, actors and objects

Ray has three primitives. A task is a stateless function call executed in some worker process: decorate a function with @ray.remote, call .remote(), and you get back an object reference immediately while the work runs elsewhere. An actor is a remote instance of a class: one long-lived process that holds state, such as a model already loaded on a GPU, and handles method calls in order. An object is an immutable value stored in the cluster's distributed object store and addressed by its reference; ray.get(ref) blocks until the value exists and fetches it.

Under the hood, each node runs a raylet, which schedules work onto local worker processes and manages that node's shared-memory object store (Plasma). The head node also runs the Global Control Store (GCS), which holds cluster metadata: which nodes exist, which actors live where, and which placement groups are reserved. Small objects travel inline with the call; large ones are written to the object store and fetched by reference, zero-copy for NumPy arrays read on the same node.

One Ray cluster running GPU workHead nodeGCS: actors, PGs, jobsWorker node A (8 GPUs)rayletlocal schedulerPlasma object storehost shared memoryActornum_gpus=1Tasknum_gpus=0.5Worker node B (8 GPUs)rayletlocal schedulerPlasma object storehost shared memoryTrain workers x8placement group bundlesDriverray.get / ray.putschedulescheduleobjects move node to nodeover the network (host memory)GPU tensors leave the GPU before they enter the object store; NCCL traffic betweentraining workers bypasses the object store and goes GPU to GPU over NVLink or the fabric.
Figure 1. Ray's moving parts on a two-node GPU cluster. Scheduling decisions use logical resources; data moves through host-memory object stores; collective traffic between training workers goes over NCCL, outside Ray.
import ray

ray.init()                      # or ray.init(address="auto") inside a cluster

@ray.remote(num_gpus=1)
class Embedder:
    def __init__(self, name):
        import torch
        from transformers import AutoModel
        self.model = AutoModel.from_pretrained(name).half().cuda().eval()

    def embed(self, input_ids):
        import torch
        with torch.inference_mode():
            x = torch.as_tensor(input_ids, device="cuda")
            return self.model(x).last_hidden_state[:, 0].float().cpu().numpy()

actors = [Embedder.remote("BAAI/bge-small-en-v1.5") for _ in range(4)]
refs = [actors[i % 4].embed.remote(batch) for i, batch in enumerate(batches)]
vectors = ray.get(refs)          # blocks until all 4 actors finish their share

Notice the shape of that code. The expensive part, loading weights onto a GPU, happens once per actor in __init__. Every later call reuses the warm model. That is the single most important GPU pattern in Ray: put state that is costly to build, such as models, CUDA contexts and compiled kernels, in actors, and use tasks for cheap stateless work.

How Ray assigns GPUs

Ray schedules on logical resources. When a node starts, Ray detects its GPUs and advertises a resource called GPU with that count. A task or actor that asks for num_gpus=1 is placed only on a node with one unit free, and the unit is subtracted until the work finishes. Before your code runs, Ray sets CUDA_VISIBLE_DEVICES in that worker process to the GPU IDs it assigned, so frameworks see the assigned device as cuda:0. ray.get_gpu_ids() returns the IDs from inside the worker.

Three consequences follow, and each one surprises people:

  • Accounting, not isolation. Ray does not stop a process from touching GPU memory it was not assigned. If code ignores CUDA_VISIBLE_DEVICES (for example, a library that resets it), two workers can land on one physical GPU.
  • Fractional GPUs share memory without limits. num_gpus=0.25 lets four actors share a GPU, which is good for small models, but Ray does not partition memory. Each actor must cap its own usage, for example with torch.cuda.set_per_process_memory_fraction or a serving engine's memory setting, or the fourth one to load will fail with out-of-memory.
  • Requesting zero GPUs can hide them. In current releases a task with no GPU request typically gets an empty CUDA_VISIBLE_DEVICES, so .cuda() fails with "CUDA is not available" even on a GPU node. Ray has signalled this default may change, so request GPUs explicitly rather than relying on either behaviour.

On mixed clusters, accelerator_type in the remote options restricts placement to nodes with a given accelerator, and custom resources let you tag nodes with properties you schedule on.

Placement groups and gang scheduling

Distributed training needs all or nothing. Eight workers that each need a GPU are useless if only six can start; the six will sit in NCCL initialisation waiting for peers that never arrive, holding GPUs another job could use. A placement group reserves a set of resource bundles atomically: either every bundle is reserved, or none is. Its strategy controls locality:

StrategyPlacementUse for
STRICT_PACKAll bundles on one node, or failTensor parallelism that needs NVLink inside one box
PACKAs few nodes as possible, best effortData parallelism where fewer hops help but are not required
SPREADAcross nodes, best effortReplicas that should not share a failure domain
STRICT_SPREADEach bundle on a different node, or failOne worker per node, for example per-node data loaders
from ray.util.placement_group import placement_group
from ray.util.scheduling_strategies import PlacementGroupSchedulingStrategy

# 8 GPUs and 10 CPUs per bundle, all on one 8-GPU node
pg = placement_group([{"GPU": 1, "CPU": 10}] * 8, strategy="STRICT_PACK")
ray.get(pg.ready(), timeout=600)            # fail loudly instead of hanging forever

workers = [
    TPWorker.options(
        num_gpus=1, num_cpus=10,
        scheduling_strategy=PlacementGroupSchedulingStrategy(
            placement_group=pg, placement_group_bundle_index=i),
    ).remote(rank=i, world_size=8)
    for i in range(8)
]

Ray Train and the vLLM integration build placement groups for you, but the same rule applies when you read their configuration: tensor-parallel groups belong inside one NVLink domain, data-parallel replicas can span nodes. Always put a timeout on pg.ready(). A group that can never be satisfied, such as nine GPUs under STRICT_PACK on 8-GPU nodes, otherwise waits silently.

The object store and GPU tensors

The object store lives in host shared memory. When an actor returns a CUDA tensor, Ray serialises it, which means copying it to the CPU first, and a consumer on another GPU copies it back. For a 4 GB activation that is two PCIe crossings plus, across nodes, a network transfer. Ray is designed for control flow and moderate data between GPU processes, not for passing large tensors between them.

So the rule is: move control through Ray, move tensors through NCCL. Training workers that need to exchange gradients form a torch.distributed process group (Ray Train sets this up) and talk GPU to GPU. Inference pipelines return small results, such as token IDs, embeddings in float16 or scores, rather than hidden states. Recent Ray releases are adding ways to keep tensors on the device between actors; treat these as version-specific and check the docs for your release before depending on them.

When the store fills, Ray spills objects to disk, which is slow on network disks. Size the store for your largest in-flight working set and point spilling at local NVMe.

Batch inference with Ray Data

Ray Data is a streaming execution engine. It reads data in blocks, runs CPU stages (decode, tokenise, resize) on CPU workers and GPU stages on a pool of GPU actors, with bounded queues between them so the GPUs are fed without loading the whole dataset into memory. The callable-class form of map_batches is the GPU pattern: the constructor runs once per actor, __call__ runs per batch.

import numpy as np
import ray

class Classifier:
    def __init__(self):
        import torch
        self.model = torch.jit.load("/models/resnet50.pt").cuda().eval().half()

    def __call__(self, batch):
        import torch
        x = torch.as_tensor(batch["image"], device="cuda").half()
        with torch.inference_mode():
            batch["label"] = self.model(x).argmax(dim=1).cpu().numpy()
        del batch["image"]                   # don't ship pixels back
        return batch

ds = (
    ray.data.read_parquet("s3://bucket/images/")
    .map_batches(decode_and_resize, num_cpus=1)           # CPU stage
    .map_batches(Classifier,
                 compute=ray.data.ActorPoolStrategy(size=16),
                 num_gpus=1, batch_size=256)               # GPU stage
)
ds.write_parquet("s3://bucket/labels/")

In current docs the actor pool is set with compute=ray.data.ActorPoolStrategy(...); the older concurrency argument is marked deprecated in recent versions, so check which one your version expects.

Worked example. Suppose 50 million images, 16 GPUs, and a model that classifies 2,000 images per second per GPU at batch 256. GPU capacity is 32,000 images per second, so the ideal run is about 26 minutes. Decoding a JPEG and resizing it costs roughly 4 ms of CPU, so feeding 32,000 per second needs about 128 busy CPU cores. On nodes with 8 GPUs and 96 cores, that is 192 cores across two nodes: enough, but only if decode is not starved by other work. If the Ray Data progress view shows the GPU stage idle and the CPU stage queue empty, you are decode-bound; add CPU nodes, which Ray can schedule the CPU stage onto, rather than more GPUs.

Distributed training with Ray Train

Ray Train wraps an ordinary PyTorch training function. You write the per-worker loop; Ray starts one worker per GPU, sets up the process group, places the workers, and handles checkpoints and restarts.

import os, tempfile, torch
import ray.train
from ray.train import ScalingConfig, RunConfig, FailureConfig, Checkpoint
from ray.train.torch import TorchTrainer, prepare_model, prepare_data_loader

def train_loop(config):
    model = prepare_model(build_model())         # moves to this worker's GPU, wraps in DDP
    opt = torch.optim.AdamW(model.parameters(), lr=config["lr"])
    loader = prepare_data_loader(build_loader())  # adds a DistributedSampler
    start = 0
    ckpt = ray.train.get_checkpoint()             # set after a restart
    if ckpt:
        with ckpt.as_directory() as d:
            state = torch.load(os.path.join(d, "state.pt"))
            model.module.load_state_dict(state["model"])
            opt.load_state_dict(state["opt"])
            start = state["epoch"] + 1
    for epoch in range(start, config["epochs"]):
        for x, y in loader:
            loss = torch.nn.functional.cross_entropy(model(x), y)
            opt.zero_grad(); loss.backward(); opt.step()
        with tempfile.TemporaryDirectory() as d:
            checkpoint = None                     # only rank 0 uploads
            if ray.train.get_context().get_world_rank() == 0:
                torch.save({"model": model.module.state_dict(),
                            "opt": opt.state_dict(), "epoch": epoch},
                           os.path.join(d, "state.pt"))
                checkpoint = Checkpoint.from_directory(d)
            ray.train.report({"loss": loss.item(), "epoch": epoch},
                             checkpoint=checkpoint)

trainer = TorchTrainer(
    train_loop,
    train_loop_config={"lr": 3e-4, "epochs": 10},
    scaling_config=ScalingConfig(num_workers=16, use_gpu=True),
    run_config=RunConfig(storage_path="s3://bucket/runs",
                         failure_config=FailureConfig(max_failures=3)),
)
result = trainer.fit()

Three details carry the reliability. storage_path must be shared storage (S3, GCS or NFS) on multi-node runs, because the replacement worker may land on a different machine. FailureConfig(max_failures=...) defaults to zero, meaning no automatic recovery. And recovery restarts the whole worker group from the latest reported checkpoint, not just the failed worker, so checkpoint frequency sets how much work a failure costs. Ray Train has a revamped "V2" implementation; the docs describe enabling it with RAY_TRAIN_V2_ENABLED=1 from Ray 2.43, so confirm which behaviour your version uses. For models that do not fit one GPU, put FSDP inside the same loop.

Failure modes

The failures below account for most Ray GPU incidents.

SymptomCauseFix
Job hangs at start, GPUs reserved, nothing runsPlacement group can never be satisfied, or NCCL waiting for a peer that was not scheduledTimeout on pg.ready(); check bundle sizes against node shapes; set NCCL timeouts
CUDA OOM on one GPU, others idleFractional num_gpus with no per-process memory cap, or a library ignoring CUDA_VISIBLE_DEVICESCap memory per actor; verify ray.get_gpu_ids() matches torch.cuda.current_device() mapping
Throughput collapses mid-runObject store full and spilling to slow diskReturn smaller results; enlarge the store; spill to local NVMe
Actor dies and job failsmax_restarts left at default for plain actors; no FailureConfig for TrainSet max_restarts and max_task_retries for idempotent actors; set max_failures for Train
Whole cluster stallsHead node lost; GCS state goneRun GCS fault tolerance with an external Redis on KubeRay, and keep the head node free of heavy work
GPUs at 30% utilisationCPU decode or a single-threaded driver loop is the bottleneckProfile stages in the Ray Data view; add CPU workers; avoid ray.get inside tight loops

A related trap is calling ray.get on each result inside a loop, which serialises the pipeline. Submit work, then consume results with ray.wait as they finish.

When Ray is worth it

Ray is not always the right tool. For a single fixed-size training job on a dedicated cluster, torchrun under Slurm or a Kubernetes training operator is simpler: fewer processes, no object store, and the same NCCL path. Ray earns its overhead when the workload is heterogeneous or dynamic: CPU preprocessing feeding GPU inference, RL loops where rollout actors and a learner exchange data, hyperparameter sweeps that start and stop trials, or one cluster serving and fine-tuning at once. On Kubernetes, KubeRay runs Ray clusters as custom resources and is the usual production deployment, with the Ray autoscaler adding and removing GPU pods.

The costs are real: another control plane to monitor, a head node that is a single point of failure without GCS fault tolerance, and abstractions that can hide where bytes go.

For further reading on this site, see PyTorch FSDP for sharding inside a Ray Train loop, GPU checkpointing for making restarts cheap, the GPU data pipeline for diagnosing input stalls, GPU hardware faults for what kills nodes, and Ray Serve for LLM systems for the online-serving side.

What to do next

  1. Run a single-node cluster with ray start --head, start one GPU actor, and print ray.get_gpu_ids() and CUDA_VISIBLE_DEVICES from inside it.
  2. Convert one batch inference script to Ray Data with a callable class, and compare GPU utilisation before and after.
  3. Put a timeout on every placement group and every ray.get in production code.
  4. Move one training job to TorchTrainer with storage_path on shared storage and max_failures set, then kill a worker node and confirm it resumes from the last checkpoint.
  5. Audit every actor that returns tensors: return CPU NumPy or small results, never large CUDA tensors.
  6. Size the object store and spill directory on GPU nodes, and alert on spill volume.
Key takeaway: Ray schedules GPUs as logical resources and sets CUDA_VISIBLE_DEVICES, but it does not isolate memory, so fractional GPUs need explicit caps. Keep costly state such as loaded models in actors, reserve multi-GPU jobs with placement groups and timeouts, move tensors over NCCL rather than the host-memory object store, and give Ray Train shared storage plus a nonzero max_failures so a lost node costs one checkpoint interval, not the run.