Ray Distributed Computing Cheat Sheet
Scale Python workloads across clusters with Ray core tasks and actors plus Ray Train, Tune, and Serve for distributed ML workflows.
Tasks and Actors
Turn a plain function into a distributed task and a class into a stateful actor.
import rayray.init() # or ray.init(address="auto") to join an existing cluster@ray.remotedef square(x): return x * xfutures = [square.remote(i) for i in range(10)]results = ray.get(futures)@ray.remoteclass Counter: def __init__(self): self.n = 0 def incr(self): self.n += 1 return self.ncounter = Counter.remote()ray.get([counter.incr.remote() for _ in range(5)]) # -> 5
Distributed Training with Ray Train
Scale a PyTorch training loop across multiple GPUs/nodes with minimal code changes.
from ray.train.torch import TorchTrainerfrom ray.train import ScalingConfigdef train_loop_per_worker(config): model = build_model() model = ray.train.torch.prepare_model(model) for epoch in range(config["epochs"]): loss = train_one_epoch(model) ray.train.report({"loss": loss})trainer = TorchTrainer( train_loop_per_worker, train_loop_config={"epochs": 10}, scaling_config=ScalingConfig(num_workers=4, use_gpu=True),)result = trainer.fit()
Hyperparameter Search with Ray Tune
Run a distributed hyperparameter sweep with an early-stopping scheduler.
from ray import tunefrom ray.tune.schedulers import ASHASchedulerdef objective(config): for step in range(20): acc = train_step(config["lr"], config["batch_size"]) tune.report({"accuracy": acc})tuner = tune.Tuner( objective, param_space={"lr": tune.loguniform(1e-4, 1e-1), "batch_size": tune.choice([16, 32, 64])}, tune_config=tune.TuneConfig(scheduler=ASHAScheduler(metric="accuracy", mode="max"), num_samples=50),)results = tuner.fit()print(results.get_best_result().config)
Serve a Model with Ray Serve
Deploy a Python class as an autoscaling HTTP inference endpoint.
from ray import serve@serve.deployment(num_replicas=2, ray_actor_options={"num_gpus": 0.5})class Predictor: def __init__(self): self.model = load_model() async def __call__(self, request): data = await request.json() return {"prediction": self.model.predict(data["input"])}serve.run(Predictor.bind(), route_prefix="/predict")
Cluster CLI Essentials
Commands for launching and managing a Ray cluster.
- ray start --head- starts the head node of a Ray cluster on the local machine
- ray start --address=<head_ip>:6379- joins a worker node to an existing cluster
- ray status- shows current cluster resource usage and node count
- ray dashboard- opens the web UI for tasks, actors, and logs
- ray.init(address="auto")- connects a script to a running cluster instead of starting one locally
Object Store: put, get, and ObjectRefs
Put large objects into Ray's shared-memory object store once and pass zero-copy references into many tasks instead of re-serializing them.
import rayimport numpy as npray.init()big_array = np.random.rand(10_000_000)ref = ray.put(big_array) # copied into the object store once@ray.remotedef sum_slice(arr_ref, start, end): return arr_ref[start:end].sum()# every task shares the same underlying memory-mapped objectfutures = [sum_slice.remote(ref, i * 1_000_000, (i + 1) * 1_000_000) for i in range(10)]total = sum(ray.get(futures))# ray.get accepts a timeout so a stuck task doesn't hang the driver forevertry: result = ray.get(ref, timeout=5.0)except ray.exceptions.GetTimeoutError: print("still running")
Placement Groups for Gang Scheduling
Reserve a bundle of resources across nodes atomically so tightly-coupled actors (e.g. parameter server + workers) land together or not at all.
from ray.util.placement_group import placement_groupfrom ray.util.scheduling_strategies import PlacementGroupSchedulingStrategypg = placement_group( bundles=[{"CPU": 4, "GPU": 1} for _ in range(4)], strategy="STRICT_PACK", # or SPREAD / PACK / STRICT_SPREAD)ray.get(pg.ready()) # blocks until the whole group is scheduledworker = Worker.options( scheduling_strategy=PlacementGroupSchedulingStrategy(placement_group=pg, placement_group_bundle_index=0)).remote()
Actor Restarts and Custom Resources
Configure automatic actor restarts and request custom (non-CPU/GPU) resources declared on cluster nodes.
@ray.remote(max_restarts=3, max_task_retries=2, resources={"TPU": 1})class FaultTolerantWorker: def __init__(self): self.state = load_checkpoint() def step(self): return train_step(self.state)worker = FaultTolerantWorker.remote()# if the actor process dies, Ray transparently restarts it (state is lost# unless you checkpoint externally) and retries in-flight tasks up to max_task_retries
Streaming ETL with Ray Data
Build a distributed, streaming data pipeline that reads, transforms, and feeds batches into training without materializing the full dataset.
import rayds = ray.data.read_parquet("s3://bucket/events/")ds = ds.map_batches(preprocess, batch_format="pandas", num_cpus=1)ds = ds.filter(lambda row: row["amount"] > 0)for batch in ds.iter_batches(batch_size=256, prefetch_batches=2): train_on_batch(batch)
Scheduling & Runtime Env Concepts
Advanced knobs for controlling where tasks run and what environment they run in.
- num_cpus=0- marks a task/actor as schedulable on any node regardless of CPU availability, useful for lightweight coordinators
- scheduling_strategy="SPREAD"- forces tasks across distinct nodes instead of packing onto one
- runtime_env={"pip": [...]}- ships a per-job Python environment to every worker without rebuilding the cluster image
- ray.remote(concurrency_groups={...})- partitions an actor's methods into separate concurrency lanes so slow calls don't block fast ones
- ray.get_runtime_context()- exposes the current task/actor ID, node ID, and namespace at runtime
- @ray.remote(num_returns="streaming")- turns a task into a generator that yields ObjectRefs incrementally instead of one final result
Set ray_actor_options={"num_gpus": 0.5} in Ray Serve deployments to pack two lightweight model replicas onto a single GPU — fractional resource requests are honored by Ray's scheduler, not just documented as a nice idea.