What’s Ray Core?#
Ray Core is a distributed computing framework for building and scaling distributed applications. It provides three essential primitives: tasks, actors, and objects. This walkthrough introduces each one with examples that show how to turn your Python functions and classes into Ray tasks and actors, and how to work with Ray objects.
Note
Ray offers an experimental API that transfers objects over Gloo, NCCL, NIXL, or your own transport, as an alternative to the default object store, which uses shared memory and gRPC. For details, see Ray Direct Transport.
Getting started#
To get started, install Ray with pip install -U ray. For other installation options, see Installing Ray.
Start by importing and initializing Ray:
import ray
ray.init()
Note
If you don’t call ray.init() explicitly, the first Ray remote API call implicitly calls ray.init() with no arguments.
Running a task#
Tasks are the simplest way to parallelize your Python functions across a Ray cluster. To create and run a task, do the following:
Decorate your function with
@ray.remoteto mark it to run remotely.Call the function with
.remote()instead of a normal function call.Use
ray.get()to retrieve the result from the returned future, which Ray calls an object ref.
The following example creates and runs a task:
# Define the square task.
@ray.remote
def square(x):
return x * x
# Launch four parallel square tasks.
futures = [square.remote(i) for i in range(4)]
# Retrieve results.
print(ray.get(futures))
# -> [0, 1, 4, 9]
Calling an actor#
Tasks are stateless. Ray actors are stateful workers that maintain their internal state between method calls. When you instantiate an actor, the following happens:
Ray starts a dedicated worker process somewhere in your cluster.
The actor’s methods run on that specific worker and can access and modify its state.
The actor executes method calls serially in the order it receives them, which preserves consistency.
The following example defines and calls a Counter actor:
# Define the Counter actor.
@ray.remote
class Counter:
def __init__(self):
self.i = 0
def get(self):
return self.i
def incr(self, value):
self.i += value
# Create a Counter actor.
c = Counter.remote()
# Submit calls to the actor. These calls run asynchronously but in
# submission order on the remote actor process.
for _ in range(10):
c.incr.remote(1)
# Retrieve final actor state.
print(ray.get(c.get.remote()))
# -> 10
The preceding example shows basic actor usage. For a fuller example that combines tasks and actors, see the Monte Carlo Pi estimation example.
Passing objects#
Ray’s distributed object store manages data across your cluster. You work with objects in Ray in three main ways:
Implicit creation: When tasks and actors return values, Ray automatically stores them in its distributed object store and returns object refs that you can retrieve later.
Explicit creation: Use
ray.put()to place objects in the store directly.Passing references: Pass object refs to other tasks and actors, which avoids unnecessary data copying and supports lazy execution.
The following example shows each technique:
import numpy as np
# Define a task that sums the values in a matrix.
@ray.remote
def sum_matrix(matrix):
return np.sum(matrix)
# Call the task with a literal argument value.
print(ray.get(sum_matrix.remote(np.ones((100, 100)))))
# -> 10000.0
# Put a large array into the object store.
matrix_ref = ray.put(np.ones((1000, 1000)))
# Call the task with the object reference as an argument.
print(ray.get(sum_matrix.remote(matrix_ref)))
# -> 1000000.0
Next steps#
Tip
To monitor your application’s performance and resource usage, see the Ray dashboard.
You can combine Ray’s primitives to express virtually any distributed computation pattern. To learn more about Ray’s key concepts, see the following user guides: