Skip to content
Back to student guides
RayMLOpsDistributed compute3 levels104 sectionsCovers Ray 2.58

The Complete Ray Guide

Scale Python and ML workloads — data, training, tuning — across a cluster with Ray. Taught at three levels — Beginner, Mid-level and Senior — each with an in-depth guide, interview prep, and practical tips.

Official docs AI-drafted · community review in progressHelp review it
15sections
22examples

This is part one of three. It covers everything you need to do real work with Ray, not a teaser. By the end you can turn an ordinary Python function into work that runs on every core of your machine, keep state in a long-lived worker, share large data without copying it, start a small cluster, submit a script to it, read the dashboard, and recognise the dozen errors that account for most beginner pain. The guide is written against Ray 2.58, the stable release at the time of writing. Mid-level and Senior take the same topics further; nothing here is thrown away.

Each section ends with a Try it task. Do them as you go. They take a few minutes each, and distributed computing only becomes intuitive once you have watched your own tasks fan out, finish out of order, and come back.

What Ray is, and the problem it solves

Ray is an open-source framework for running Python code on many processors at once, on one machine or across a cluster of machines, while keeping the programming model close to ordinary Python. You mark a function or a class as something that may run elsewhere, call it in a slightly different way, and Ray decides where it runs, moves the data it needs, and hands you the result.

To see why that is useful, start with what a Python program does by default: it runs one thread of work at a time. Python's global interpreter lock means that even if you start several threads, only one of them executes Python bytecode at any instant. If your laptop has ten cores, a plain Python loop uses roughly one of them. For a quick script this does not matter. For the workloads that machine learning produces it matters enormously: tuning forty hyperparameter combinations, resizing two million images, running a model over a night's worth of documents, or simulating ten thousand game episodes for reinforcement learning. All of those are collections of independent pieces of work, and the natural thing is to run the pieces side by side.

Before Ray, you had a few options, none of them comfortable. Python's built-in multiprocessing module spreads work over cores on one machine, but it stops at the machine boundary and gives you little help sharing large data. Cluster frameworks such as Apache Spark (see the Spark guide) scale across machines, but they ask you to express your program in their data-processing vocabulary, which fits table-shaped transformations well and fits "run this Python class on a GPU and keep it alive" poorly. Message queues and job schedulers work across machines but leave you writing the plumbing by hand: serialising arguments, tracking which worker has what, retrying failures. Ray's pitch is to give you the plumbing, with a small API that feels like Python.

PYTHON FUNCTIONwhat you wrote
→
@ray.remotemark it distributable
→
.remote() CALLSmany run in parallel
→
ray.get()collect the results

The diagram is the whole core idea, and the rest of this guide fills it in.

Ray has two layers, and keeping them apart will save you confusion. Ray Core is the low-level runtime: tasks, actors and objects, which you will meet in the next section. On top of Core sit five libraries aimed at machine-learning work: Ray Data for loading and transforming data and running batch inference, Ray Train for distributed model training, Ray Tune for hyperparameter search, Ray Serve for online model serving (it has its own guide, Ray Serve), and RLlib for reinforcement learning. The libraries are built from Core primitives; understanding Core means you understand what the libraries are doing when something goes wrong.

What people use it for:

⚡

Parallel Python

Turn a slow loop over independent items into work spread over every core you own, with a few changed lines.

🧠

Distributed training

Train a model across several GPUs or machines and have failed workers handled for you.

🎛️

Hyperparameter search

Run many training trials at once and stop the unpromising ones early.

🚀

Model serving

Run a model behind an HTTP endpoint with several replicas that scale with traffic.

You need little to follow along: Python 3, a terminal, and a machine with a few cores. Nothing here requires a GPU or a cloud account.

Try it
  1. Find a loop in one of your own scripts where each iteration is independent of the others.
  2. Write down roughly how long one iteration takes and how many there are.
  3. Multiply, then divide by the number of cores on your machine.
a number much smaller than the one you started with. That gap is what Ray is for, and you will close part of it in the first project below.

The mental model: tasks, actors and objects

Ray Core has three nouns. If you can explain these to a colleague, you understand more of Ray than most people who have copied a tutorial.

A task is a function that runs somewhere else. You decorate an ordinary Python function with @ray.remote, and from then on you can call it with .remote() instead of directly. A task is stateless: it receives arguments, computes, returns a value, and is forgotten. Because tasks share nothing, Ray can run many at once and run each wherever there is room.

An actor is a class whose instance lives in its own worker process. You decorate a class with @ray.remote, create an instance with .remote(), and call its methods with .remote() too. Unlike a task, an actor remembers: its attributes persist between calls. Actors are the answer to "I need something that loads a large model once and then answers many requests" or "I need one place that keeps a running count". By default an actor handles one method call at a time, in the order received, so you do not need locks around its state.

An object is a value stored in Ray's distributed object store. Every .remote() call returns immediately with an ObjectRef, which is a handle to a result that may not exist yet. Think of it as a receipt, or in programming terms a future or promise. You turn the receipt into the actual value by calling ray.get(ref), which waits until the work is finished. Objects are immutable: once stored, a value never changes, which is what lets Ray share them freely between processes.

Surrounding these three nouns is the machinery that runs them.

head node
GCSGlobal Control Service. Cluster metadata, the actor directory, placement groups
DashboardWeb UI and Jobs API, port 8265 by default
RayletThe per-node scheduler and object-store manager
each worker node
RayletReceives work, starts worker processes
Worker processesRun your tasks and host your actors
Object storeShared memory holding objects, spilling to disk when full

A cluster is one head node plus any number of worker nodes. The head runs the Global Control Service (GCS), which stores cluster metadata, and the dashboard. Every node, the head included, runs a raylet, which decides what runs on that node and manages its slice of the object store. The driver is simply your own program, the process where you call ray.init(). Workers are the processes that actually execute tasks and host actors. On your laptop all of this lives on one machine; Ray starts it for you when you call ray.init(), which is why you can learn the whole model without a cluster.

One more concept deserves attention early: resources. Each node advertises logical quantities of CPUs, GPUs and memory, and each task or actor declares how much it needs (one CPU by default for a task). The scheduler places work where the declared resources are free. These are scheduling bookkeeping, not enforcement. If you say a task needs one CPU and it spins up eight threads, Ray will not stop it. Ray uses the numbers to avoid overbooking, so the honest numbers are your side of the bargain.

Three nouns, one rule Tasks are stateless functions, actors are stateful classes, objects are immutable values. When you are unsure which to use, ask whether the work needs to remember anything between calls. If not, it is a task. If so, it is an actor.
Try it
  1. For each case, say task or actor: resizing one image; a service that holds a loaded language model; adding two numbers; a shared counter of how many jobs finished.
  2. For each, say what you would pass around as an object.
resizing and adding are tasks, the model holder and the counter are actors, and the images, model weights or numbers are the objects.

Installing Ray and checking the setup

Ray is installed with pip. Always do it inside a virtual environment so that Ray's dependencies stay separate from your system Python.

BASH
python3 -m venv .venv
source .venv/bin/activate
pip install -U "ray[default]"

The bare pip install -U "ray" gives you Core only. The default extra adds the dashboard, the Jobs system and the cluster launcher, which you will want from the second half of this guide onwards. If you plan to use the libraries, extras exist for them, such as ray[data,train,tune,serve]. To pin an exact release, which is good practice because Ray ships a new minor version every few weeks, write pip install "ray[default]==2.58.0".

On Windows, the commands are slightly different:

BASH
py -m venv .venv
.venv\Scripts\activate
pip install -U "ray[default]"

Support differs by platform. Linux on x86_64 and aarch64 and macOS on Apple silicon are the well-supported targets. Windows is marked beta in the documentation: multi-node clusters on Windows are untested and file opening is slower. If you are on Windows, running Ray inside WSL2 is a practical way to get Linux behaviour. Check the installation page for the range of Python versions the current release supports before you create your virtual environment, because the supported range moves between releases.

If you prefer conda, the conda-forge channel carries a ray-default package, and Docker images are published as rayproject/ray. When running Ray inside a Docker container, give the container shared memory, for example --shm-size=2gb. Without it, Ray warns that the object store is using /tmp instead of /dev/shm and runs noticeably slower.

Now verify the install two ways, once from the shell and once from Python:

BASH
ray --version
python -c "import ray; ray.init(); print(ray.cluster_resources())"

The second command starts a private one-machine Ray instance, prints the resources it found, and exits. The output is a dictionary along these lines, with numbers that depend on your hardware:

TEXT
{'CPU': 10.0, 'memory': 14000000000.0, 'object_store_memory': 6000000000.0, 'node:127.0.0.1': 1.0}

Read it as a list of what Ray believes it has to spend: CPU is the number of logical cores it detected, memory and object_store_memory are byte counts, and the node: entry is a special per-node resource Ray uses internally. If you have a supported NVIDIA GPU and the right drivers, a GPU key appears too. The exact numbers will differ on your machine; what matters is that the CPU count matches your core count, because that number caps how many one-CPU tasks run at once.

Pin the version in anything you share A tutorial written for an older Ray may import modules that have since moved or been removed. When you copy a snippet and it fails on import, check the date of the tutorial before you debug your code. In a project, record ray==2.58.0 (or whatever you tested) in your requirements file.
Try it
  1. Create a virtual environment and install ray[default].
  2. Run the verification one-liner and compare the CPU value with your actual core count.
a printed dictionary whose CPU count equals your logical cores. If the import fails, the virtual environment is probably not activated.

Your first project: from a slow loop to parallel tasks

Create a file called first.py. The program below pretends to do slow work, first serially and then with Ray, and times both. The time.sleep stands in for real computation such as a feature calculation or an API call.

first.py
import time
import ray

ray.init()

def slow_square(x):
    time.sleep(1)          # pretend this is real work
    return x * x

@ray.remote
def slow_square_remote(x):
    time.sleep(1)
    return x * x

# Serial: ordinary Python, one at a time
start = time.time()
serial = [slow_square(i) for i in range(8)]
print("serial  :", serial, f"{time.time() - start:.1f}s")

# Parallel: submit everything, then collect
start = time.time()
refs = [slow_square_remote.remote(i) for i in range(8)]
parallel = ray.get(refs)
print("parallel:", parallel, f"{time.time() - start:.1f}s")

ray.shutdown()

Run it with python first.py. On a machine with at least eight cores, the output looks like this, with the serial run near eight seconds and the parallel one a little over one:

TEXT
serial  : [0, 1, 4, 9, 16, 25, 36, 49] 8.0s
parallel: [0, 1, 4, 9, 16, 25, 36, 49] 1.3s

The extra third of a second is Ray's start-up cost: launching worker processes the first time they are needed. If your machine has four cores, expect the parallel run to take about two seconds, since only four one-CPU tasks run at a time and eight tasks need two rounds. That is the resource system doing its job.

Walk through the changes, because there are only four. First, ray.init() starts the runtime. Second, @ray.remote on the function turns it into a remote function. Third, calling it as slow_square_remote.remote(i) returns an ObjectRef at once, without waiting, so the list comprehension finishes in milliseconds having merely submitted eight tasks. Fourth, ray.get(refs) accepts a list of references and blocks until all of them are ready, returning the values in the same order as the list.

That ordering is worth saying out loud. The tasks may finish in any order, but ray.get on a list always returns results in the order you submitted them, so you can pair results with inputs without extra bookkeeping.

If you call the remote function directly, slow_square_remote(3), Ray raises an error telling you to use .remote() instead. That is the most common first mistake and it is a helpful one: the error is the framework preventing you from accidentally running something locally that you meant to distribute.

The trap that makes Ray slower than a loop Writing ray.get(f.remote(i)) inside the loop submits one task, waits for it, then submits the next. You have paid Ray's overhead and received no parallelism. Always submit everything first, collecting the references in a list, and call ray.get once on the whole list.

Serial by accident

  • [ray.get(f.remote(i)) for i in range(8)]
  • Waits after every submission
  • Takes as long as a plain loop, plus overhead

Parallel

  • ray.get([f.remote(i) for i in range(8)])
  • Submits all eight, waits once
  • Takes about as long as the slowest task per round
Try it
  1. Run first.py and note both timings.
  2. Change range(8) to range(32) and predict the parallel time before running it.
  3. Then deliberately write the trap form inside the loop and time it.
the parallel time grows in steps of one second per round of your core count, and the trap form comes back to the serial time.

Tasks in depth

Now that the basic loop works, the details of tasks are what separate a script that happens to run from one you understand.

Tasks can call tasks. Inside a remote function you may call other remote functions and even pass their ObjectRefs along. If you pass an ObjectRef as an argument to another task, Ray waits for it to be ready and substitutes the real value before your function runs. This lets you build pipelines without ever calling ray.get in the middle:

pipeline.py
import ray
ray.init()

@ray.remote
def load(n):
    return list(range(n))

@ray.remote
def total(numbers):
    return sum(numbers)

data_ref = load.remote(1000)
result_ref = total.remote(data_ref)   # pass the ref, not the data
print(ray.get(result_ref))            # 499500

Here the list produced by load never travels to your driver at all if the two tasks land on the same node; Ray resolves the reference where it is needed. That matters when the intermediate values are large.

Getting results as they finish. ray.get on a list waits for every task. When the tasks take uneven time and you would like to act on each result as soon as it arrives, use ray.wait. It takes a list of references and returns two lists: those that are ready and those that are not.

wait_demo.py
import time, random
import ray
ray.init()

@ray.remote
def work(i):
    time.sleep(random.random() * 3)
    return i

pending = [work.remote(i) for i in range(6)]
while pending:
    ready, pending = ray.wait(pending, num_returns=1)
    print("finished:", ray.get(ready[0]))

Each loop turn prints whichever task finished first, so the numbers appear out of order. ray.wait also accepts a timeout in seconds, which returns whatever is ready when time runs out. This is the right tool for processing a stream of results, or for giving up on stragglers.

Declaring what a task needs. By default a task asks for one CPU. You can say otherwise in the decorator or per call:

resources.py
import ray
ray.init()

@ray.remote(num_cpus=2)
def heavy(x):
    return x * 2

@ray.remote
def light(x):
    return x + 1

print(ray.get(heavy.remote(5)))                        # uses the decorator's 2 CPUs
print(ray.get(light.options(num_cpus=0.5).remote(5)))  # override for this call

num_cpus, num_gpus and memory are the three you will use most. Fractional values are allowed, so num_gpus=0.5 lets two tasks share one GPU. Remember that these are promises to the scheduler. A task declared num_gpus=1 is simply given a GPU assignment; with PyTorch or TensorFlow you still have to put your tensors on that device. A task that needs more of a resource than any single node has can never run, and Ray reports it as infeasible; the errors section shows how that looks.

Failures. If your task raises an exception, Ray catches it, stores it, and re-raises it on the driver when you call ray.get, wrapped in a RayTaskError that includes the remote traceback. This means a bug in a remote function shows up in the place you are looking, with the line number from the worker. Separately, if a worker process dies, for example because the machine ran out of memory, Ray retries the task automatically a few times by default. You can control this with max_retries in the decorator. Retry on application exceptions is off by default, and you turn it on with retry_exceptions=True.

Idempotence is your responsibility. Because tasks may be retried, a task should be safe to run twice. Appending a row to a database inside a task, then having the worker crash after the insert but before returning, produces a duplicate on retry. Design tasks as "compute a value from inputs", and write side effects deliberately.

Try it
  1. Run wait_demo.py three times and notice the finish order changes.
  2. Write a task that raises ValueError("bad input") for a negative number, call it with -1 and read the traceback that ray.get raises.
  3. Rerun resources.py after changing heavy to num_cpus=64 on a laptop and watch what Ray prints.
shuffled finish order, a traceback that shows your own ValueError line, and a warning that the request cannot be scheduled, with the call hanging rather than failing.

Actors: remembering things between calls

Tasks forget everything. Sometimes that is exactly wrong. Loading a 2 GB model for each of a thousand predictions would be absurd; you want to load it once and keep it. That is what an actor is for.

counter.py
import ray
ray.init()

@ray.remote
class Counter:
    def __init__(self):
        self.n = 0

    def inc(self):
        self.n += 1
        return self.n

    def value(self):
        return self.n

c = Counter.remote()                       # starts a worker process, runs __init__ there
refs = [c.inc.remote() for _ in range(3)]
print(ray.get(refs))                       # [1, 2, 3]
print(ray.get(c.value.remote()))           # 3

Three things happen that are worth spelling out. Counter.remote() creates a dedicated worker process and runs __init__ inside it; this takes a moment, which is why you create actors up front, not per request. c.inc.remote() sends a message to that process and returns an ObjectRef for the result. And because the actor handles one call at a time in arrival order, the three increments return [1, 2, 3] without any locking.

The state lives in the actor's process, not in your driver. You cannot read c.n directly; attribute access on the handle does not reach across processes. You must define a method that returns it, as value does above, and call it with .remote().

Actors are the right shape for several common jobs:

  • Holding a loaded model so that each call only pays for inference.
  • Holding a connection or client, such as a database session, that is expensive to create and unsafe to pickle (serialise) across processes.
  • Accumulating results, such as counters, metrics or a shared best-score tracker.
  • Wrapping a stateful service, which is exactly how Ray Serve builds its replicas.

Actors also take resource declarations, and they hold their resources for as long as they live. This is the cause of a classic puzzle: you create four actors that each ask for one CPU on a four-core machine, and then your tasks never run, because every CPU is claimed by an idle actor. The fix is to size actors deliberately or give them num_cpus=0 when they are only coordinating.

Several actors can run side by side, which is how you get parallelism with state: create a handful of actors that each hold a model copy, and spread requests across them.

pool.py
import ray
ray.init()

@ray.remote
class Scorer:
    def __init__(self, scale):
        self.scale = scale             # imagine loading a model here

    def score(self, x):
        return x * self.scale

scorers = [Scorer.remote(10) for _ in range(3)]
refs = [scorers[i % 3].score.remote(i) for i in range(9)]
print(ray.get(refs))                   # [0, 10, 20, 30, 40, 50, 60, 70, 80]

When you are finished with an actor, it goes away automatically once all handles to it are out of scope and the driver exits. To stop one explicitly, call ray.kill(actor_handle). You can also give an actor a name, Counter.options(name="counter").remote(), and fetch it from another script with ray.get_actor("counter"). A named actor with lifetime="detached" outlives the script that created it, which is useful and also a way to leak resources, so use it with intent.

An actor is one process, and one lane Because an actor runs one method at a time, a slow method blocks every other call to that actor. If you route all traffic to a single actor, you have rebuilt a single-threaded server. Use several actors, or keep methods short.
Try it
  1. Run counter.py, then add a reset method and call it between increments.
  2. Create two Counter actors and increment them independently to confirm they hold separate state.
two counters that never interfere. Each actor handle points at its own process with its own self.

Objects and the object store

Every value that flows between Ray tasks and actors lives in the object store, a region of shared memory on each node. You normally never touch it directly, but understanding it explains the speed, the memory errors and one habit that makes programs much faster.

When a task returns a value, Ray writes it into the object store of the node where the task ran and gives you an ObjectRef. When another task needs it, Ray copies it across the network if the task is on another node, or lets the task read it straight from shared memory if it is on the same node. For NumPy arrays this reading is zero-copy: the array in your worker points into the shared memory rather than copying it, which is why passing large arrays between tasks on one machine is cheap.

You can place a value in the store yourself with ray.put, which returns a reference. The reason to do this is that a big argument passed directly to many tasks is serialised again for every call, whereas a reference is tiny:

put_demo.py
import numpy as np
import ray
ray.init()

big = np.random.rand(10_000_000)          # about 80 MB

@ray.remote
def head_sum(arr, k):
    return arr[:k].sum()

big_ref = ray.put(big)                    # store once
refs = [head_sum.remote(big_ref, k) for k in (10, 100, 1000)]
print(ray.get(refs))

Pass big_ref and the array is stored once and shared; pass big and every call serialises 80 MB again. Ray will usually warn you when a remote function captures a very large value in its closure, which is the same mistake in disguise.

The store has a fixed size per node. When it fills up, Ray spills older objects to disk, which keeps things working but slows them down, and if the store is full and nothing can be spilled you receive an ObjectStoreFullError. Objects are also reference-counted: when no ObjectRef points at a value any more, it is freed. This is why holding on to a giant list of references you no longer need keeps memory busy. Delete the list, or let it go out of scope, to release the data.

Two exceptions to know by name. ObjectLostError means an object disappeared, for example because the node that held it died and Ray could not rebuild it. OwnerDiedError means the process that created the reference, its owner, has died, so Ray no longer knows how to find the value. Both are rare on a laptop and common in clusters with unreliable machines; the senior material covers lineage reconstruction.

When in doubt, pass references If the same large object goes to many tasks, ray.put it once. If one task's output feeds another, pass the ObjectRef, not the result of ray.get. Pulling a large result to the driver and sending it back out is the most common source of needless copying.
Try it
  1. Run put_demo.py. Then edit it to pass big directly and time both versions over 20 calls.
  2. Call ray.get(big_ref) and check whether the result equals big.
the reference version is clearly faster for many calls, and the round-trip returns an identical array.

The cluster: starting nodes, the dashboard, and the CLI

Until now ray.init() quietly created a private one-machine instance that died with your script. Real work happens on a cluster that outlives scripts. A cluster is started with the ray command line tool, and you can build one entirely on your laptop to learn the shape.

Start a head node:

BASH
ray start --head

The command prints the address of the cluster, a command to add more nodes, a snippet to connect from Python, and the URL of the dashboard, which by default is http://127.0.0.1:8265. To join another machine as a worker you run ray start --address=<head_ip>:6379 on it, where 6379 is the default port of the Global Control Service. That is the whole procedure: one head, then each worker pointed at it. The machines must be able to reach each other over the network, and every node must run the same Ray version and the same Python version; a mismatch is refused on connect.

Now connect your script to the running cluster instead of creating a private one:

connect.py
import ray

ray.init(address="auto")      # find the cluster started on this machine
print(ray.cluster_resources())

address="auto" looks for a running Ray instance on this machine. If none exists, you get a ConnectionError, which is the signal to run ray start --head first. You can also set the RAY_ADDRESS environment variable so that both ray.init() and the CLI find the cluster without any argument.

Open the dashboard in a browser. Its pages list nodes, jobs, actors and tasks with their state, CPU and memory use, and logs. During your experiments it is the fastest answer to "what is Ray doing right now?". The same information is available in the terminal through the State API commands:

BASH
ray status            # resources in use and available, plus autoscaler state
ray list actors       # every actor and whether it is alive
ray list tasks        # recent tasks and their states
ray list nodes        # the nodes in the cluster
ray summary tasks     # tasks grouped by function name and state
ray logs              # browse and fetch logs
ray memory            # which objects are using the object store
ray stop              # stop the Ray processes on this machine

ray status deserves a habit of its own. It shows the usage as used / total for each resource, and pending demands that cannot be met. When a program seems stuck, this single command usually tells you whether Ray is waiting for resources that do not exist.

Ray has no login screen By default, anyone who can reach the dashboard or the cluster's ports can run arbitrary code on your machines. The documentation is explicit that Ray allows any client to run arbitrary code. Never expose port 8265 or 6379 to the internet. Keep the cluster on a private network, and read the senior guide before running one for a team. Token authentication exists from 2.52, as an extra layer, not a substitute for network isolation.

On a laptop the head and a worker can both run locally for practice. After ray start --head, running ray start --address=127.0.0.1:6379 in a second terminal adds a second node on the same machine, which is a perfectly valid way to see two nodes listed in the dashboard. Finish with ray stop in each terminal.

On managed platforms, the head and workers are started for you. On Kubernetes the standard route is the KubeRay operator, which has its own place in the Kubernetes guide's ecosystem, and the commercial managed platform from Ray's creators is called Anyscale. You do not need either to complete this guide.

Try it
  1. Run ray start --head and open the dashboard URL it prints.
  2. Run connect.py, then ray status and ray list nodes.
  3. Run ray stop when you are done.
one node in the dashboard, a resource table in ray status, and clean shutdown. If ray start reports that Ray is already running, run ray stop first.

Submitting work as a job, and controlling the environment

Running python script.py on your laptop and connecting with ray.init(address=...) works while you are developing. For anything repeated, Ray offers Ray Jobs: you submit a script to the cluster, the cluster runs it, and you can check its status and logs afterwards, even after your terminal is closed. Jobs go through the dashboard's address, port 8265, so install ray[default] to have them.

Save a small script:

job_script.py
import ray

ray.init()

@ray.remote
def square(x):
    return x * x

print(ray.get([square.remote(i) for i in range(10)]))

Submit it from the directory that contains it:

BASH
ray job submit --address http://127.0.0.1:8265 --working-dir . -- python job_script.py

Everything after the double dash is the command the cluster will run. --working-dir . uploads the current directory so that the cluster can find your files, which matters when the cluster is on other machines that do not have your code. The command streams the job's output and prints a job identifier. Afterwards you can inspect it:

BASH
ray job status <job_id>
ray job logs <job_id>
ray job stop <job_id>

Jobs solve a real problem. If your driver runs on a laptop that sleeps, the work dies with it; a job lives on the cluster.

The same submission can be made from Python through JobSubmissionClient, which is the way a scheduler or a web application starts Ray work, but the command line is enough to start.

Submitting code to remote machines raises an obvious question: how do the right libraries get there? The answer is runtime_env, a description of the environment that a job, task or actor needs. You can set it for the whole program in ray.init:

runtime_env_demo.py
import ray

ray.init(
    runtime_env={
        "pip": ["requests"],
        "env_vars": {"GREETING": "hello"},
        "working_dir": ".",
    }
)

@ray.remote
def read_env():
    import os
    return os.environ.get("GREETING")

print(ray.get(read_env.remote()))   # hello

pip lists packages to install for the workers, env_vars sets environment variables, and working_dir ships a folder of your code. Ray builds the environment on the workers on demand and caches it. This is far more convenient than logging into each machine to install packages, and it is why mismatched dependencies are less common with Ray than with hand-managed servers. It is not free: building a large environment takes time on first use, so for big dependencies like deep-learning frameworks you normally bake them into the machine image or container image instead.

Do not put secrets in runtime_env Values in env_vars end up in job configuration that others can view in the dashboard and logs. Pass credentials through your platform's secret mechanism, not through committed code.
Try it
  1. Start a head node, then submit job_script.py with ray job submit.
  2. Fetch its logs with ray job logs after it finishes.
  3. Find the job in the dashboard's Jobs page.
the list of ten squares in the logs, a job marked succeeded, and the same record visible in the browser.

The libraries on top of Core

You now know the foundation. This section is a guided tour of the five libraries, enough to recognise what each is for and to run a small example, not a complete treatment. All of them are installed with extras, such as pip install -U "ray[data,train,tune,serve]".

Ray Data loads, transforms and writes datasets that are bigger than one machine's memory, and is most often used for batch inference and preprocessing. A dataset is split into blocks that are processed in a streaming fashion, so memory stays bounded. Here is a small example that needs no files:

data_demo.py
import ray

ds = ray.data.range(1000)                       # a dataset with an "id" column

def double(batch):
    batch["id"] = batch["id"] * 2
    return batch

ds = ds.map_batches(double)
print(ds.take(3))

map_batches runs your function over batches of rows, which is faster than row by row, and is the central operation for running a model over a dataset: you pass a class that loads the model once and use concurrency to say how many copies run. Real datasets come from ray.data.read_parquet("s3://...") and friends, and results are written with write_parquet. Some older Ray Data functions are deprecated in recent releases, such as Dataset.zip, so if a tutorial uses them, check the release notes.

Ray Train runs a training function across several workers, which can be GPUs or machines, and takes care of starting the processes, wiring them together and recovering from failures. You write a function that trains a model for one worker, and hand it to a trainer such as TorchTrainer with a scaling configuration saying how many workers and whether they use GPUs. Inside your function you report metrics and save checkpoints through Ray Train's APIs. Ray Train has gone through a major redesign; many older tutorials use an earlier interface, so always read the docs for the version you installed.

Ray Tune runs hyperparameter search. You describe a search space, such as learning rates to try, and Tune launches many trials in parallel, each a Ray task or actor, and can stop weak trials early. It integrates with experiment trackers such as MLflow. The API centres on tune.Tuner, which takes your training function, a param_space and a tune.TuneConfig, and .fit() runs the search. The older tune.run function appears in many tutorials and is the previous style.

Ray Serve serves models online. You define a deployment, a class or function that Ray runs as several replicas, and Serve handles the HTTP routing and scaling. A deployment looks like this:

serve_demo.py
from ray import serve

@serve.deployment(num_replicas=2)
class Hello:
    async def __call__(self, request):
        return {"message": "hello from Ray Serve"}

app = Hello.bind()
serve.run(app, route_prefix="/")

After serve.run, an HTTP server on port 8000 answers requests at /. The Serve guide, Ray Serve, covers deployment configuration, autoscaling and production use in full. For serving large language models, see also vLLM, which Ray Serve integrates with.

RLlib is the reinforcement learning library. It uses Core's actors to run many simulated environments in parallel and collect experience for training. It is a specialist tool, so most newcomers can skip it until they have a reinforcement learning problem.

How do you choose? The question to ask is what your work looks like:

If your work is... Reach for
Independent pieces of Python Core tasks
Something that must hold state Core actors
Preprocessing or batch inference over a big dataset Ray Data
One model trained over many GPUs Ray Train
Trying many configurations of one model Ray Tune
An HTTP endpoint for a model Ray Serve

Because all five sit on Core, they can be combined: a Tune run can launch Train jobs, and a Serve deployment can call a Data pipeline. They also coexist with the rest of the stack. Ray can read from Parquet and Delta Lake tables (Delta Lake), and training runs started from Ray report to trackers such as MLflow.

Try it
  1. Install the extras with pip install -U "ray[data,serve]" and run data_demo.py.
  2. Run serve_demo.py in a script that sleeps afterwards (add import time; time.sleep(600)), then in another terminal request curl http://127.0.0.1:8000/.
a few rows with doubled ids, and a JSON message from the Serve endpoint. Stop the script with Ctrl+C when finished.

Configuration and resources you will actually touch

Ray has many settings, but a beginner needs only a handful, and most are arguments to ray.init or ray start.

ray.init() takes num_cpus and num_gpus to override what Ray detected, which is handy for testing. Setting num_cpus=4 on a ten-core laptop makes Ray behave as if it had four, so you can observe how your program behaves on a smaller machine. It also takes address to connect to an existing cluster, runtime_env as seen earlier, and namespace, which scopes named actors so that two jobs do not collide on a name. A namespace is for organisation, not security.

The same ideas apply on the command line: ray start --head --num-cpus=4 limits the head, and --dashboard-host controls which interface the dashboard listens on. The default listens on localhost only, which is the safe choice; changing it to 0.0.0.0 makes the unauthenticated dashboard reachable from the network, so do it only behind a firewall.

The environment variables you will meet first are RAY_ADDRESS for the default cluster address and RAY_DEDUP_LOGS, which controls the collapsing of repeated log lines from many workers into one. Treat any other RAY_ variable you find in a forum post as unfamiliar until you have checked the docs for your version, because many are internal tuning knobs that change between releases.

Logs live on each node in a session directory under /tmp/ray, with the newest session linked as session_latest. Output from print inside tasks and actors is streamed to the driver by default, prefixed with the function name and process ID, which is how you see (square pid=1234) in front of lines. When you debug, remember that those lines come from worker processes: a print in a task is not printed in your own process.

For memory, the two settings that matter early are the amount of object store memory and the per-task memory declaration. Ray has a memory monitor that kills workers when the node approaches full memory, so the system preserves itself rather than letting the operating system kill something at random. When that happens you see an out-of-memory error, covered next.

Develop small, scale later Use ray.init(num_cpus=2) while writing code, to confirm that a program still works on small resources and to make scheduling behaviour easy to follow. Remove the argument when you run for real. Ray programs scale up because the same code runs unchanged on a bigger cluster, so correctness testing on a small machine transfers.
Try it
  1. Run first.py with ray.init(num_cpus=2) and watch the parallel time.
  2. Add print inside the remote function and observe the prefix on each line.
the parallel run takes about four seconds instead of one, because eight one-second tasks need four rounds on two CPUs, and each printed line is tagged with a function name and process id.

Common errors and how to read them

Most Ray errors fall into a few families. Learn the family and you can guess the cause before reading the message in full. Messages below are shown as they appear in current documentation where known; wording shifts between releases, so match on the exception class.

RayTaskError wrapping your exception. Not a Ray problem. It means the code in your task raised something, and the traceback beneath shows the line in your function. Read from the bottom: the last lines are your real error. Fix the function.

Calling a remote function directly. The error says remote functions cannot be called directly and must be called with .remote(). Replace f(x) with f.remote(x), and Counter() with Counter.remote().

Out-of-memory kills. An OutOfMemoryError saying a task was killed because the node was running low on memory means Ray's memory monitor killed a worker that was using too much. Reduce how many tasks run at once, declare memory= for heavy tasks so Ray schedules fewer per node, or pass references rather than copies. Use a larger machine if the work genuinely needs the memory.

ObjectStoreFullError. The object store ran out of room and spilling could not free enough. Hold fewer references at once, process results in a streaming fashion with ray.wait, and delete references you no longer need.

RayActorError. An actor's process died, so calls to it fail. Causes include a crash in your method, running out of memory, or a node loss. Look at the actor's logs. If you want automatic recovery, configure max_restarts on the actor, covered in the mid-level guide.

WorkerCrashedError. The worker running a task died unexpectedly, again typically from out-of-memory or a native library crashing. Ray retries tasks a few times, so the error appears when the retries are exhausted.

GetTimeoutError. You called ray.get(ref, timeout=...) and the deadline passed. Not a failure of the task itself, which may still be running.

A program that hangs with a warning about resources. If a task or actor asks for more than the cluster has, Ray cannot schedule it and prints a warning that the request cannot be scheduled right now. The call waits indefinitely. Run ray status in another terminal, compare the pending demand with the totals, and either reduce the request or add capacity. A common cause is actors holding all CPUs, as noted earlier.

Cannot connect to a cluster. ray.init(address="auto") raises a ConnectionError when there is no running head. Start one with ray start --head, or drop address to let Ray create a local instance. If you connect to a remote head and it hangs, a firewall is probably blocking the GCS port, 6379, or the ports Ray picks for its other services; open the ports between nodes.

Version mismatch. A worker or client started with a different Ray or Python version than the head is refused, with a message naming both sets of versions. Reinstall to align them. This is the most common failure when students join a second machine to the cluster.

Pickling errors. Ray serialises your functions and arguments with cloudpickle. Objects holding locks, open files, sockets or database connections cannot be serialised, and you see errors such as TypeError: cannot pickle '_thread.lock' object. Create such objects inside the task or actor, not outside it, and pass only plain data.

Pydantic version trouble. Ray 2.56 dropped support for Pydantic version 1. If a Ray import fails with a Pydantic complaint, upgrade to Pydantic 2.

An old tutorial's import fails. Ray has reorganised its libraries over time, for example retiring the ray.air namespace and redesigning Train. When ModuleNotFoundError points at a Ray module, check the documentation for your installed version before changing anything else.

Try it
  1. Cause four errors on purpose: call a remote function without .remote(); raise an exception in a task; call ray.init(address="auto") with no cluster running; capture a threading.Lock() in a task.
  2. For each, write one sentence naming the family and the fix.
four distinct messages, each of which you can now classify without searching.

Putting it all together

Here is one small project that uses almost everything above. It scores a batch of short texts with a pretend model. A Scorer actor holds the model, so loading happens once per actor; a task function splits the text into chunks; a shared lookup table is stored once with ray.put; and results are collected with ray.wait as they arrive.

project.py
import time
import ray

ray.init(num_cpus=4)

STOPWORDS = {"the", "a", "of", "and", "to"}
stop_ref = ray.put(STOPWORDS)            # shared once, passed by reference

@ray.remote
def clean(text, stop):
    words = [w for w in text.lower().split() if w not in stop]
    return " ".join(words)

@ray.remote
class Scorer:
    def __init__(self):
        time.sleep(1)                    # pretend to load a model
        self.calls = 0

    def score(self, cleaned):
        self.calls += 1
        return len(cleaned.split())      # pretend score: word count

    def count(self):
        return self.calls

texts = [
    "The cost of the cloud and the cost of GPUs",
    "A guide to the basics of Ray",
    "Tasks and actors and objects",
    "To scale a Python program",
    "The object store holds the data",
    "Jobs run on the cluster",
]

scorers = [Scorer.remote() for _ in range(2)]
cleaned_refs = [clean.remote(t, stop_ref) for t in texts]
score_refs = [scorers[i % 2].score.remote(ref) for i, ref in enumerate(cleaned_refs)]

pending = list(score_refs)
results = {}
while pending:
    ready, pending = ray.wait(pending, num_returns=1)
    idx = score_refs.index(ready[0])
    results[texts[idx]] = ray.get(ready[0])

for text, score in results.items():
    print(score, "|", text)

print("calls per scorer:", ray.get([s.count.remote() for s in scorers]))
ray.shutdown()

Walk through it as a checklist of ideas. The stopword set is stored once with ray.put and the reference goes to every clean call. The clean tasks run in parallel and return references. Those references are passed straight to the actors' score methods, so cleaned text never visits the driver. Two Scorer actors share the load, and each pays its start-up cost once. ray.wait lets you handle each score as it finishes. The last line asks each actor how many calls it served, and the answer, three and three, shows that state really did live in the actors.

Run it and expect six lines of the form 5 | The cost of the cloud and the cost of GPUs in the order they finished, then calls per scorer: [3, 3]. Because of the one-second pretend model load, the whole run takes a couple of seconds.

Now take it one step further as an exercise. Start a head with ray start --head, change ray.init(num_cpus=4) to ray.init(address="auto"), and submit it as a job with ray job submit --working-dir . -- python project.py. Open the dashboard and find the two actors and the tasks. This is the complete everyday Ray workflow: write locally, test small, submit to a cluster, observe.

Try it
  1. Run project.py and confirm the calls per scorer.
  2. Change to three scorers and predict the call counts, then verify.
  3. Run it as a job and find its actors in ray list actors.
each of three scorers serving two calls, and actors listed as alive while the job is running.

What you can now do, and what comes next

You can explain Ray's three nouns and the head, worker and raylet structure behind them. You can parallelise independent work with tasks, avoid the loop-with-ray.get trap, process results as they arrive with ray.wait, hold state in actors, share large values by reference, start a small cluster, open its dashboard, query it with ray status and ray list, submit a script as a job with a runtime_env, and read the common errors by family. You also know what each of the five libraries is for and which one fits which kind of work.

Keep three habits. Submit all work before you collect any, pass references instead of data, and assume every task can run twice. Keep two warnings in mind: Ray has no authentication by default, so never expose its ports publicly, and Ray releases often, so pin the version and read the release notes before upgrading. If your data or your employer's rules require that it stay in a particular country, the cluster's machines must be placed in a cloud region that satisfies it; for teams in the Gulf and Egypt that means choosing in-region infrastructure for the cluster, not just for storage.

Where to go next depends on what you want to do. The mid-level guide covers placement groups, fault tolerance in depth, Ray Data and Train in practice and the Jobs API. The senior guide covers running Ray on Kubernetes, high availability for the head node, security and multi-tenancy. Related guides in this catalogue: Ray Serve for online inference, Kubernetes as the usual production home, Docker for building the images your workers run, Apache Spark for a different take on distributed data, and MLflow for tracking what your Ray jobs produce.

Sources