How the Fully Async RL Pipeline in Miles Decouples Rollout and Training Workers
The Fully Async RL pipeline in Miles decouples rollout and training workers by running environment rollouts in a persistent background Ray actor with a queue-based data buffer, while the training loop fetches batches asynchronously and overlaps the next rollout preparation with current training steps.
Miles (radixark/miles) is an open-source RL framework that implements a fully-asynchronous reinforcement learning architecture. Unlike synchronous pipelines where generation and training block each other, Miles separates these concerns through three interconnected components: a dedicated rollout executor, an inference controller, and independent training models. This design eliminates global barriers and maximizes GPU utilization.
Core Components of the Fully Async RL Pipeline
Rollout Executor with Persistent Background Worker
The rollout_executor runs as a Ray actor that continuously generates rollout data—token sequences, rewards, and other trajectory information—without pausing for the training loop.
Key implementation details from miles/rollout/fully_async_rollout.py:
- The executor owns a
FullyAsyncDataBufferthat queues generated rollout groups - This buffer smooths production-consumption mismatches, allowing the trainer to fetch data at its own pace
- Metrics including queue length, staleness, and dropped groups are exposed for monitoring under
rollout/fully_async/...
# From fully_async_rollout.py — persistent worker with data buffer
class FullyAsyncRolloutFn:
def __init__(self):
self.data_buffer = FullyAsyncDataBuffer() # Decouples producer/consumer
# ... worker runs continuously in background Ray actor
Inference Controller for Coordination
The inference_controller orchestrates between the trainer and executor without forcing synchronization points:
| Function | Purpose |
|---|---|
prepare_rollout() |
Prepares the model state for the next rollout generation |
update_weights() primitive |
Pushes policy updates asynchronously to the executor |
The controller exposes a thin RPC surface. The trainer requests rollouts via rollout_executor.get.remote() without blocking on weight synchronization.
Training Models with Async Weight Updates
The actor_model and critic_model perform RL training steps on rollout data fetched from the executor. They receive policy weights through the same asynchronous channel used by the executor, ensuring approximate synchronization without global barriers.
Pipeline Execution Flow in train_async.py
The driver script train_async.py implements the following sequence:
1. Component Initialization (Lines 35-38)
inference_controller, rollout_executor, num_rollout = await create_rollout_components(args)
actor_model, critic_model = await create_training_models(
args, inference_controller, rollout_executor
)
2. Initial Weight Synchronization (Line 53)
await update_weights(actor_model, rollout_executor) # Non-blocking push to executor
3. Main Async Loop with Overlapped Operations (Lines 76-89)
# Eagerly start first rollout preparation
rollout_future = await eager_create_task(prepare_and_generate(start_rollout_id))
for rid in range(start_rollout_id, args.num_rollout):
# Await the rollout produced by background worker
rollout_data = await rollout_future
# Immediately kick off next rollout (overlaps with training)
if rid + 1 < args.num_rollout:
rollout_future = await eager_create_task(prepare_and_generate(rid + 1))
# Train on current rollout while next one generates
await actor_model.train(rid, rollout_data)
4. Periodic Weight Updates (Line 119)
if (rid + 1) % args.update_weights_interval == 0:
await update_weights(actor_model, rollout_executor, rollout_id=rid)
The eager_create_task primitive is critical: it ensures the next rollout preparation begins before training completes, achieving true pipelining.
Key Decoupling Mechanisms
| Mechanism | How Decoupling Works |
|---|---|
| Ray actor isolation | Rollout executor runs in separate process/thread from trainer |
FullyAsyncDataBuffer |
Queue-based buffer absorbs rate mismatches between producer and consumer |
Asynchronous get.remote() |
Training loop pulls batches without blocking on rollout generation |
update_weights channel |
Weight pushes happen independently; no global synchronization barrier |
| Eager task creation | Next rollout starts early, overlapping generation with training |
Complete Implementation Example
import asyncio
from miles.rollout.fully_async_rollout import FullyAsyncRolloutFn
from miles.train_async import (
create_rollout_components,
create_training_models,
update_weights,
eager_create_task
)
async def fully_async_training_loop(args):
# Build async components
inference_controller, rollout_executor, num_rollout = await create_rollout_components(args)
# Create models with references to executor for data fetching
actor_model, critic_model = await create_training_models(
args, inference_controller, rollout_executor
)
# Initial policy sync to executor
await update_weights(actor_model, rollout_executor)
# Overlapped rollout generation and training
rollout_future = await eager_create_task(
inference_controller.prepare_rollout(start_rollout_id)
and rollout_executor.get.remote(start_rollout_id)
)
for rid in range(start_rollout_id, args.num_rollout):
# Consume rollout from buffer
rollout_data = await rollout_future
# Start next rollout immediately (non-blocking)
if rid + 1 < args.num_rollout:
rollout_future = await eager_create_task(
inference_controller.prepare_rollout(rid + 1)
and rollout_executor.get.remote(rid + 1)
)
# Train while next rollout generates in background
await actor_model.train(rid, rollout_data)
await critic_model.train(rid, rollout_data)
# Periodic weight push without stopping executor
if (rid + 1) % args.update_weights_interval == 0:
await update_weights(actor_model, rollout_executor, rollout_id=rid)
return actor_model
Relevant Source Files
| File Path | Description |
|---|---|
miles/rollout/fully_async_rollout.py |
FullyAsyncRolloutFn class with persistent worker and data buffer logic |
miles/rollout/fully_async_data_buffer.py |
Queue implementation tracking staleness, queue length, dropped groups |
train_async.py |
Driver orchestrating inference controller, executor, and training models |
miles/utils/arguments.py |
--fully-async flag parsing and rollout implementation selection |
tests/fast/rollout/test_fully_async_rollout.py |
Unit tests for buffer metrics and staleness handling |
Summary
- Rollout executor runs as persistent Ray actor with dedicated
FullyAsyncDataBufferqueue - Inference controller coordinates without blocking through
prepare_rolloutandupdate_weightsprimitives - Training loop fetches batches asynchronously via
get.remote()and overlaps next rollout with current training - Weight updates flow through async channel, avoiding global barriers while keeping policy roughly synchronized
eager_create_taskenables true pipelining by starting rollouts before training completes
Frequently Asked Questions
How does Miles prevent the training loop from starving when rollout generation is slow?
The FullyAsyncDataBuffer maintains a queue of generated rollouts. The training loop consumes from this buffer asynchronously via await rollout_executor.get.remote(). If generation lags, the trainer simply waits; if generation runs ahead, excess rollouts queue up. The buffer exposes metrics (rollout/fully_async/queue_length) to monitor this balance.
What happens to rollouts generated with stale policy weights?
Miles tracks staleness explicitly in fully_async_data_buffer.py. Each rollout records its generation policy version, and the buffer calculates the gap between consumed rollout version and current trainer version. While some staleness is inherent to async training, the update_weights_interval parameter bounds it—weights push periodically to refresh the executor's policy.
Can the Fully Async pipeline run on a single GPU?
Yes, though the design targets multi-GPU scenarios. On single GPU, the Ray actor executor and trainer share the device with time-slicing. The prepare_rollout call in the inference controller handles resource allocation. The async structure still provides latency hiding benefits even without true parallelism.
How does eager_create_task differ from standard asyncio.create_task?
eager_create_task (implemented in Miles utilities) ensures the task begins execution immediately rather than waiting for the next event loop iteration. This minimizes the gap between requesting and starting the next rollout, tightening the overlap between generation and training.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →