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 FullyAsyncDataBuffer that 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 FullyAsyncDataBuffer queue
  • Inference controller coordinates without blocking through prepare_rollout and update_weights primitives
  • 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_task enables 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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →