asyncio and concurrency for actors#

A single actor process can run concurrent threads.

Ray offers two types of concurrency within an actor:

Keep in mind that Python’s Global Interpreter Lock (GIL) runs only one thread of Python code at a time.

As a result, parallelizing pure Python code doesn’t give you true parallelism. If you call NumPy, Cython, TensorFlow, or PyTorch code, these libraries release the GIL when they call into C or C++ functions.

Note

Neither the threaded actor nor the async actor model bypasses the GIL.

asyncio for actors#

Python 3.5 and later support writing concurrent code with the async/await syntax. Ray integrates natively with asyncio, so you can use Ray alongside popular async frameworks such as aiohttp and aioredis.

import ray
import asyncio

@ray.remote
class AsyncActor:
    def __init__(self, expected_num_tasks: int):
        self._event = asyncio.Event()
        self._curr_num_tasks = 0
        self._expected_num_tasks = expected_num_tasks

    # Multiple invocations of this method can run concurrently on the same event loop.
    async def run_concurrent(self):
        self._curr_num_tasks += 1
        if self._curr_num_tasks == self._expected_num_tasks:
            print("All coroutines are executing concurrently, unblocking.")
            self._event.set()
        else:
            print("Waiting for other coroutines to start.")

        await self._event.wait()
        print("All coroutines ran concurrently.")

actor = AsyncActor.remote(4)
refs = [actor.run_concurrent.remote() for _ in range(4)]

# Fetch results using regular `ray.get`.
ray.get(refs)

# Fetch results using `asyncio` APIs.
async def get_async():
    return await asyncio.gather(*refs)
asyncio.run(get_async())
(AsyncActor pid=9064) Waiting for other coroutines to start.
(AsyncActor pid=9064) Waiting for other coroutines to start.
(AsyncActor pid=9064) Waiting for other coroutines to start.
(AsyncActor pid=9064) All coroutines are executing concurrently, unblocking.
(AsyncActor pid=9064) All coroutines ran concurrently.
(AsyncActor pid=9064) All coroutines ran concurrently.
(AsyncActor pid=9064) All coroutines ran concurrently.
(AsyncActor pid=9064) All coroutines ran concurrently.
...

Object refs as asyncio.Future objects#

You can use object refs as asyncio.Future objects, which means you can await Ray futures in existing concurrent applications.

Instead of:

import ray

@ray.remote
def some_task():
    return 1

ray.get(some_task.remote())
ray.wait([some_task.remote()])

you can await the object ref in Python 3.9 and 3.10:

import ray
import asyncio

@ray.remote
def some_task():
    return 1

async def await_obj_ref():
    await some_task.remote()
    await asyncio.wait([some_task.remote()])

asyncio.run(await_obj_ref())

or await the Future object directly in Python 3.11 and later:

import asyncio

async def convert_to_asyncio_future():
    ref = some_task.remote()
    fut: asyncio.Future = asyncio.wrap_future(ref.future())
    print(await fut)
asyncio.run(convert_to_asyncio_future())
1

See the asyncio documentation for more asyncio patterns, including timeouts and asyncio.gather.

Object refs as concurrent.futures.Future objects#

You can also wrap object refs in concurrent.futures.Future objects to work with existing concurrent.futures APIs:

import concurrent

refs = [some_task.remote() for _ in range(4)]
futs = [ref.future() for ref in refs]
for fut in concurrent.futures.as_completed(futs):
    assert fut.done()
    print(fut.result())
1
1
1
1

Define an async actor#

Ray automatically detects whether an actor supports async calls from its async method definitions.

import ray
import asyncio


@ray.remote
class AsyncActor:
    def __init__(self, expected_num_tasks: int):
        self._event = asyncio.Event()
        self._curr_num_tasks = 0
        self._expected_num_tasks = expected_num_tasks

    async def run_task(self):
        print("Started task")
        self._curr_num_tasks += 1
        if self._curr_num_tasks == self._expected_num_tasks:
            self._event.set()
        else:
            # Yield the event loop for multiple coroutines to run concurrently.
            await self._event.wait()

        print("Finished task")

actor = AsyncActor.remote(5)
# All 5 tasks will start at once and run concurrently.
ray.get([actor.run_task.remote() for _ in range(5)])
(AsyncActor pid=3456) Started task
(AsyncActor pid=3456) Started task
(AsyncActor pid=3456) Started task
(AsyncActor pid=3456) Started task
(AsyncActor pid=3456) Started task
(AsyncActor pid=3456) Finished task
(AsyncActor pid=3456) Finished task
(AsyncActor pid=3456) Finished task
(AsyncActor pid=3456) Finished task
(AsyncActor pid=3456) Finished task

Ray runs all of the methods inside a single Python event loop. Don’t run blocking ray.get or ray.wait calls inside an async actor method, because ray.get blocks the event loop.

An async actor runs only one task at any point in time, though Ray can multiplex tasks on it. An async actor has only one thread. If you want a thread pool, see Threaded actors.

Set concurrency in async actors#

Use the max_concurrency flag to set how many “concurrent” tasks run at once. By default, 1000 tasks can run concurrently.

import asyncio
import ray

@ray.remote
class AsyncActor:
    def __init__(self, batch_size: int):
        self._event = asyncio.Event()
        self._curr_tasks = 0
        self._batch_size = batch_size

    async def run_task(self):
        print("Started task")
        self._curr_tasks += 1
        if self._curr_tasks == self._batch_size:
            self._event.set()
        else:
            await self._event.wait()
            self._event.clear()
            self._curr_tasks = 0

        print("Finished task")

actor = AsyncActor.options(max_concurrency=2).remote(2)

# Only 2 tasks will run concurrently.
# Once 2 finish, the next 2 should run.
ray.get([actor.run_task.remote() for _ in range(8)])
(AsyncActor pid=5859) Started task
(AsyncActor pid=5859) Started task
(AsyncActor pid=5859) Finished task
(AsyncActor pid=5859) Finished task
(AsyncActor pid=5859) Started task
(AsyncActor pid=5859) Started task
(AsyncActor pid=5859) Finished task
(AsyncActor pid=5859) Finished task
(AsyncActor pid=5859) Started task
(AsyncActor pid=5859) Started task
(AsyncActor pid=5859) Finished task
(AsyncActor pid=5859) Finished task
(AsyncActor pid=5859) Started task
(AsyncActor pid=5859) Started task
(AsyncActor pid=5859) Finished task
(AsyncActor pid=5859) Finished task

Threaded actors#

Sometimes asyncio isn’t the right fit for your actor. For example, you might have a method that performs a computation-heavy task and blocks the event loop without giving up control through await. This hurts the performance of an async actor, because async actors can only run one task at a time and rely on await to switch context.

Instead, use the max_concurrency actor option without any async methods to get threaded concurrency, similar to a thread pool.

Warning

If an actor definition has at least one async def method, Ray recognizes the actor as an async actor instead of a threaded actor.

@ray.remote
class ThreadedActor:
    def task_1(self): print("I'm running in a thread!")
    def task_2(self): print("I'm running in another thread!")

a = ThreadedActor.options(max_concurrency=2).remote()
ray.get([a.task_1.remote(), a.task_2.remote()])
(ThreadedActor pid=4822) I'm running in a thread!
(ThreadedActor pid=4822) I'm running in another thread!

Each invocation of a threaded actor runs in a thread pool. The max_concurrency value limits the size of the thread pool.

asyncio for remote tasks#

Ray doesn’t support asyncio for remote tasks. The following snippet fails:

@ray.remote
async def f():
    pass

Instead, wrap the async function in a wrapper that runs the task synchronously:

async def f():
    pass

@ray.remote
def wrapper():
    import asyncio
    asyncio.run(f())