Fully Async Policy · 技术分享

Harbor × VeRL 全异步 RL 训练架构详解

从启动脚本到权重同步:一条样本如何走过 Rollouter、Agent Loop、Harbor 沙箱、轨迹 Proxy、MessageQueue,再进入 Trainer。

源文档:Verl+harbor全异步RL分享.md
主线示例:fully_async_2nodes.sh(1 节点 8 GPU Rollouter + 1 节点 8 GPU Trainer)

怎么读这篇

第 1 节讲为什么拆成两组 GPU以及四个旋钮;第 2 节是并发漏斗,读完就能回答「64 个并发卡在哪」;第 3–5 节是 Harbor 适配与 Token-in-Token-out;第 6–7 节是训练步与 Partial Rollout。附录是总图、文件索引和术语。

本文档以 fully_async_2nodes.sh 的启动流程为例,系统性地介绍从 VeRL 全异步训推框架、Agent Loop 抽象、Harbor SWE Agent Loop 适配、轨迹代理 Proxy、到训练与权重同步的完整链路。

SWE-Lego RL Trainer 整体架构:沙箱、进程内 Proxy、负载均衡、Rollout Engine 与 Trainer
SWE-Lego RL Trainer 总览。点击可打开原图。来源:SWE-Lego-RL-Trainer-framework-v2.pdf
上 · Environment Harbor 沙箱:K8s / Docker / 云主机,对应第 4 节
中 · In-Process Proxy 进程内 aiohttp 假 vLLM,对应第 5 节
中 · Global Load Balancer 粘性路由 + 最少在途请求,对应第 2.5 节
下 · Rollout ↔ Trainer 样本缓冲、权重同步、Partial Rollout,对应第 1 / 6 / 7 节

1VeRL 全异步训推模式介绍

1.1 背景与动机

在传统的 Colocate(训推一体)架构中,Rollout 和 Training 共享同一组 GPU 资源。这种模式下存在严重的资源浪费问题:

全异步训练(Fully Async Policy)通过资源隔离 + 并行生成与训练 + 流式样本传输彻底解耦了 Rollouter 和 Trainer,实现了 2.35x-2.67x 的训练吞吐提升。

1.2 核心架构四组件

01 · 生产者

Rollouter

独占一组 GPU,持续流式生成样本,以单样本为最小传输单元放入 MessageQueue。生成速度受 staleness 约束。

02 · 通道

MessageQueue

Ray Actor 异步队列,deque + asyncio.Condition 实现生产者-消费者,逐条传递样本。

03 · 消费者

Trainer

独占另一组 GPU,从队列取样本,攒够 require_batches × ppo_mini_batch_size 后执行一轮 PPO。

04 · 同步

ParameterSynchronizer

基于 NCCL,参考 checkpoint-engine,把 Trainer 新权重高效广播回 Rollouter。

┌─────────────────────────────────────────────────────────────────────┐
│                       Fully Async Policy                             │
│                                                                      │
│  ┌──────────┐    sample by sample   ┌──────────────┐                │
│  │ Rollouter │ ──────────────────> │ MessageQueue │                 │
│  │ (N GPUs) │                      │  (Ray Actor) │                 │
│  └──────────┘                       └──────┬───────┘                │
│       ▲                                    │ sample by sample        │
│       │ NCCL                               ▼                         │
│  ┌────┴──────────────┐              ┌──────────────┐                │
│  │ ParameterSynchronizer │<─────────│   Trainer    │                 │
│  │ (checkpoint-engine)  │          │  (M GPUs)    │                 │
│  └──────────────────────┘           └──────────────┘                │
└─────────────────────────────────────────────────────────────────────┘
  1. Rollouter:独占一组 GPU(如 8 卡),持续流式生成样本,以单样本为最小传输单元放入 MessageQueue。生成速度受 staleness(陈旧度)参数控制。
  2. MessageQueue:基于 Ray Actor 的异步消息队列,使用 deque + asyncio.Condition 实现生产者-消费者模式。
  3. Trainer:独占另一组 GPU(如 8 卡),从 MessageQueue 逐个取出样本,攒够 require_batches * ppo_mini_batch_size 个样本后执行一轮 PPO 训练。
  4. ParameterSynchronizer:基于 NCCL 通信原语,参考 checkpoint-engine,实现 Trainer 到 Rollouter 的高效参数同步。

1.3 收益来源

Colocate · 串行
|------ rollout 1(含长尾)------|-- train --|------ rollout 2 ------|
                  ↑
        绝大多数样本早已跑完,
        整批在等最慢的 1~2 条
Fully Async · 重叠
Rollouter  |--- 持续生成 ---|--- 持续生成 ---|--- 持续生成 ---|
Trainer             |---- train 1 ----|---- train 2 ----|
           ═══════> 两条时间线重叠,端到端大幅缩短

1.4 四种运行模式

Mode a

On Policy Pipeline

sync step
= 1
staleness
0
partial

完全同步:每次训完立即同步参数。

Mode b

Stream Off Policy

sync step
> 1
staleness
0
partial

流式同步:多步本地训练后才同步参数。

Mode c

Async Stream(stale)

sync step
≥ 1
staleness
> 0
partial
False

允许陈旧样本;同步时等待活跃任务完成。

Mode d

Async Stream(partial)

sync step
≥ 1
staleness
> 0
partial
True

允许陈旧样本;同步时中断并恢复进行中的 rollout。

模式 trigger_parameter_sync_step staleness_threshold partial_rollout 特点
a. On Policy Pipeline =1 0 - 完全同步,每次训完立即同步参数
b. Stream Off Policy Pipeline >1 0 - 流式同步,多步训练后才同步参数
c. Async Stream (stale) >=1 >0 False 允许陈旧样本,同步时等待活跃任务完成
d. Async Stream (partial) >=1 >0 True 允许陈旧样本,同步时中断并恢复进行中的 rollout
四种全异步 RL 模式对比图
四种运行模式的官方对比图。若内网图无法加载,以上方四张卡片和表格为准。

1.5 关键参数说明

async_training.trigger_parameter_sync_step

Trainer 执行多少次本地更新后,才与 Rollouter 做一次参数同步。一次本地更新 = 取到 require_batches × ppo_mini_batch_size 个样本。

两次同步之间 Trainer 处理的样本数 =
trigger_parameter_sync_step × require_batches × ppo_mini_batch_size

与 Colocate 公平比速度时,应设为 data.train_batch_size / (require_batches × ppo_mini_batch_size)

async_training.staleness_threshold

允许使用的过期样本最大比例。num_staleness_sample 是上一窗口多生成、带进本窗口的过期样本数。Rollouter 慢时 Trainer 会更早触发同步,实际产不出满额 rollout_num。速度足够快且阈值设为 1,大致等价于 one-step-off。为避免过期样本伤精度,建议 小于 1

阈值 = 0(同步):rollout_num = trigger × require_batches × ppo_mini_batch_size
阈值 > 0(异步):rollout_num = (1 + threshold) × trigger × require_batches × ppo_mini_batch_size − num_staleness_sample

async_training.partial_rollout

仅在 staleness_threshold > 0 时真正生效。为 True 时,权重同步会中断进行中的生成,同步后再从中断点续写。

async_training.require_batches

流式训练里通常设为 1:攒够一个 ppo_mini_batch_size 就训。实测一次分发太少会因数据顺序导致训练不稳、响应变长,所以用这个旋钮控制每次参训样本量。

actor_rollout_ref.actor.use_rollout_log_probs=True

PPO / GRPO / DAPO 的重要性采样要求 old_log_prob 与 rollout 当时的参数和 token 对齐。全异步默认由 Rollout 侧计算,而不是 Trainer 重算。

algorithm.rollout_correction.bypass_mode

默认 True:直接用 rollout log prob。训练后期指标或响应长度不稳时,可关掉它,改用训练引擎重算 log_prob,并启用 Rollout Importance Sampling。在 mode d + Partial Rollout 下,这近似 AReaL 的 Decoupled PPO。

1.6 快速启动示例

fully_async_2nodes.sh 为例,启动命令的核心是:

python -m verl.experimental.fully_async_policy.fully_async_main \
    --config-name='harbor_verl_fully_async_fsdp.yaml' \
    --config-path="$REPO_ROOT/src/verl_patch/config" \
    actor_rollout_ref.actor.ppo_mini_batch_size=64 \
    actor_rollout_ref.rollout.n=8 \
    async_training.staleness_threshold=1.0 \
    async_training.trigger_parameter_sync_step=1 \
    async_training.partial_rollout=True \
    async_training.require_batches=1 \
    rollout.nnodes=1 rollout.n_gpus_per_node=8 \
    trainer.nnodes=1 trainer.n_gpus_per_node=8 \
    # ... 其他训练超参数

这条命令启动了两个 Ray Actor——FullyAsyncTrainer(1 节点 8 GPU)和 FullyAsyncRollouter(1 节点 8 GPU),通过 MessageQueue 进行异步通信。


2VeRL Agent Loop 抽象与并发控制

2.1 Agent Loop 类层次结构

AgentLoopBase  (抽象基类)
├── SingleTurnAgentLoop   (单轮 agent,常规的单轮reasoning任务比如AIME,MATH500都用这个)
├── DiffusionSingleTurnAgentLoop  (扩散模型 agent)
├── ToolAgentLoop          (多轮工具调用 agent)
│   └── (AsyncPartialToolAgentLoop 能力通过 FullyAsyncLLMServerManager 注入)
└── BuiltinSWEAgentLoop    (适配harbor的 SWE-bench agent,直接继承 AgentLoopBase)

2.2 AgentLoopBase 核心抽象

agent_loop.py 中定义了核心数据结构和基础类:

# AgentLoopOutput 是单条轨迹的完整输出
class AgentLoopOutput:
    prompt_ids: list[int]         # prompt 的 token id 序列
    response_ids: list[int]       # response 的 token id 序列
    response_mask: list[int]      # 1=LLM 生成的 token, 0=tool/observation token
    response_logprobs: list[float] # 每个 response token 的 log 概率
    routed_experts: Any           # MoE 路由专家选择
    reward_score: float           # 奖励分数
    num_turns: int                # 对话轮次数
    metrics: dict                 # 性能指标
    extra_fields: dict            # 额外元数据

AgentLoopBase 的核心方法:

2.3 ToolAgentLoop 状态机 (无执行环境)

ToolAgentLoop 实现了标准的多轮工具调用循环:

                    ┌─────────┐
                    │ PENDING │  (初始状态:应用 chat template 得到 prompt_ids)
                    └────┬────┘
                         ▼
    ┌──────────────────────────────────────┐
    │              GENERATING              │ (调用 server_manager.generate())
    │  - 累积 response tokens (mask=1)     │
    │  - 提取 tool calls                   │
    │  - 检查终止条件                       │
    └───────┬──────────────┬───────────────┘
            │              │
    无 tool call      有 tool call
            │              │
            ▼              ▼
      ┌──────────┐  ┌──────────────────┐
      │TERMINATED│  │ PROCESSING_TOOLS │  (并行执行 tool calls, mask=0)
      └──────────┘  └────────┬─────────┘
                              │
                              ▼
                    回到 GENERATING (下一轮)

状态由 AgentData 对象承载,贯穿整个状态机,包含 messages、token ids、masks、logprobs 等全部轨迹状态。

2.4 并发控制全链路分析

全异步模式下,从 Rollouter 发起一个样本到最终在 Harbor 沙箱中完成执行,需要经过多层并发限制。下面以 fully_async_2nodes.sh 的配置为例(n_gpus_rollout=8, gen_tp=4, num_workers=128),逐层分析全链路的并发瓶颈。

2.4.1 配置回顾与关键数值计算

# 来自 fully_async_rollouter.py set_max_required_samples()

# 1. 最大陈旧样本数:一次参数同步窗口内,Rollouter 最多生成的样本总量
max_required_samples = required_samples * (staleness_threshold + 1) * trigger_parameter_sync_step
                    = (64 * 1) * (1.0 + 1) * 1
                    = 128

# 2. vLLM Server Replica 数量
num_replicas = n_gpus_rollout / tensor_model_parallel_size
             = 8 / 4 = 2   # 每个 replica 占用 4 张 GPU

# 3. 最大并发样本数
max_concurrent_samples = min(len(server_handles) * 32, max_required_samples)
                       = min(2 * 32, 128)
                       = min(64, 128) = 64

# 4. 最大队列容量
max_queue_size = max_required_samples = 128

2.4.2 第一层:pending_queue(数据加载 → 任务提交的缓冲区)

位置FullyAsyncRollouter 内的 asyncio.Queue(maxsize=128)

[_feed_samples()]                          [_processor_worker()]
从 dataloader 逐条读取 batch_dict         从 pending_queue 取出样本
→ 包装为 RolloutSample                    → 创建 asyncio.Task 提交处理
→ await pending_queue.put(sample)         → 加入 active_tasks 集合

并发限制:Queue 有界(maxsize=128),feed 协程在 Queue 满时自动阻塞
位置:fully_async_rollouter.py:427-458 (_feed_samples)
      fully_async_rollouter.py:461-544 (_processor_worker)

作用:防止 dataloader 无限读取数据撑爆内存。feed 协程与 processor 协程以"生产者-消费者"模式解耦。

实际并发效果:Queue 中的样本尚未开始被处理,它们只是等待处理的数据缓冲。

2.4.3 第二层:active_tasks + max_concurrent_samples(任务提交的硬限制)

位置FullyAsyncRollouter._processor_worker() 中的提交循环

# fully_async_rollouter.py:527-534
# 在提交新任务前,检查当前活跃任务数是否超限
while len(self.active_tasks) >= self.max_concurrent_samples:
    done_tasks, self.active_tasks = await asyncio.wait(
        self.active_tasks, return_when=asyncio.FIRST_COMPLETED
    )
    for task in done_tasks:
        await task

# 未超限时,创建一个新 asyncio.Task
task = safe_create_task(
    self._process_single_sample_streaming(rollout_sample),
    name=rollout_sample.sample_id,
    task_set=self.active_tasks,
)

作用:这是 Rollouter 侧对"同时在途样本数"的唯一硬限制(max_concurrent_samples = 64)。当 64 个样本都处于处理中(从提交给 async_rollout_manager 到结果写入 MessageQueue 之间)时,processor_worker 会阻塞直到至少一个任务完成。

2.4.4 第三层:AgentLoopWorker 池(Worker 级别的串行化)

位置FullyAsyncAgentLoopManager.generate_sequences_single() 中的 asyncio.Queue Worker 池

# fully_async_policy/agent_loop/agent_loop.py:206-230
async def generate_sequences_single(self, prompts: DataProto) -> DataProto:
    queue = await self._init_worker_queue()  # 预填充 128 个 Worker
    single_samples = list(prompts.chunk(len(prompts)))

    async def process_task(sample: DataProto) -> DataProto:
        worker = await queue.get()     # 从池中获取一个空闲 Worker
        try:
            result = await worker.generate_sequences.remote(sample)
            # ↑ 在 Worker 内部: asyncio.gather(*tasks)
            # tasks 来自 batch 中的每条样本(这里每 batch 就是一条样本)
            return result
        finally:
            await queue.put(worker)   # 归还 Worker 到池中

    outputs = await asyncio.gather(*[process_task(sample) for sample in single_samples])

作用num_workers = 128FullyAsyncAgentLoopWorker(每个是一个 Ray Actor)构成 Worker 池。generate_sequences_single 被调用时,一个 batch(可能包含多条样本,在 n_resp_per_prompt=8 时就是 8 条)被拆分为单条样本,每条样本通过 queue.get() 获取一个空闲 Worker 来执行。

实际并发效果:最多有 128 个 AgentLoopWorker 同时在处理样本。但由于上游 max_concurrent_samples = 64,实际同时在途的样本不会超过 64 个,所以这一层在当前配置下不是瓶颈。

2.4.5 第四层:单个 AgentLoopWorker 内部的并发(batch 内样本并行)

位置AgentLoopWorker.generate_sequences()

# agent_loop.py:601-611
tasks = []
for i in range(len(batch)):   # batch 中有多少条样本
    tasks.append(asyncio.create_task(
        self._run_agent_loop(sampling_params, trajectory_info[i], ...)
    ))
outputs = await asyncio.gather(*tasks)

作用:一个 AgentLoopWorker 收到一个 batch 后,对 batch 内的每条样本(每个 agent_nameBuiltinSWEAgentLoop.run())都创建一个 asyncio.Taskgather 并行执行。所以一个 Worker 可同时运行多条样本的 BuiltinSWEAgentLoop.run()

2.4.6 第五层:Harbor 沙箱侧(K8s Pod / Docker 容器的资源上限)

Harbor 侧的并发没有额外的软件层限制——K8s Pod 的创建数量仅受 K8s 集群资源(CPU/内存配额)限制。这意味着:当 128 个样本同时进入 Harbor Trial.run() 时,会同时创建多达 128 个 K8s Pod。这通常需要 K8s 集群有足够的节点资源,否则 Pod 会排队等待调度。

2.4.7 第六层:vLLM Server(GPU 推理的物理上限)

每个 LLM 推理请求最终通过 FullyAsyncLLMServerManager.generate()AsyncLLMServerManager.generate() 发送到 vLLM Server replica。

Rollout 侧 8 GPU, gen_tp=4 → 2 个 vLLM Server Replica
每个 Replica 是一个独立的 vLLM 实例,可处理多个并发请求

GlobalRequestLoadBalancer 通过 sticky session + least-inflight-requests
将请求分配到具体的 vLLM Server

实际限制:2 个 vLLM Server Replica 承载 64 个并发样本中的所有 LLM 推理请求。由于每个样本是多轮交互(LLM 推理 → 工具执行 → LLM 推理 → ...),并非所有请求同时到达 vLLM,实际的 vLLM 并发压力取决于各样本的节奏交错情况。

2.4.8 全链路并发总结

pending_queue · 128
数据缓冲,满则阻塞 dataloader
active_tasks · 64
Rollouter 唯一软件硬限制
AgentLoopWorker 池 · 128
当前配置下不是瓶颈
Harbor / K8s Pod · ≤64
受集群 CPU / 内存配额约束
vLLM · 2 replicas(dp2 tp4)
GPU 推理物理上限,通常是最紧的一层
Dataloader
  │ 逐条读取
  ▼
pending_queue (maxsize=128, asyncio.Queue)
  │ 取出提交
  ▼
active_tasks (max_concurrent=64, asyncio.Task 集合)
  │ 每个 task 调用:
  │   async_rollout_manager.generate_sequences_single()
  ▼
AgentLoopWorker 池 (num_workers=128, asyncio.Queue)
  │ 每个 Worker 处理 1 条样本
  │ 内部 asyncio.gather 并行执行 batch 内样本
  ▼
BuiltinSWEAgentLoop.run()
  │ 启动 Proxy → Harbor Trial.run()
  ▼
┌─────────────────────────────────────┐
│            Harbor 沙箱               │
│  K8s Pod / Docker Container          │
│  ├─ Agent 启动                       │
│  └─ Agent 循环:                     │
│      ├─ HTTP → Proxy → vLLM (2 replicas)│  ← GPU 推理的物理上限
│      └─ 工具执行 (bash/edit/git)     │
└─────────────────────────────────────┘
  │ 结果序列化,写入 MessageQueue
  ▼
MessageQueue (Ray Actor, deque)
  │
  ▼
FullyAsyncTrainer 消费

并发瓶颈排序(从紧到松):

位置 并发上限 限制类型 备注
vLLM Server Replicas dp2tp4(GPU 物理资源) 硬件 8 GPU / gen_tp=4,所有 LLM 推理的物理上限
active_tasks 64 软件硬限制 min(2*32, 128),Rollouter 侧的唯一并发硬限制
pending_queue 128 软件缓冲 缓冲 128 个待处理样本,超过则阻塞 feed
AgentLoopWorker 池 128 软件容量 当前低于 active_tasks × 8 的理论上限,因此不是瓶颈
K8s 集群资源 取决于集群 基础设施 64 个 Pod 需要足够的 CPU/内存

2.5 GlobalRequestLoadBalancer:粘性会话与最少请求数负载均衡

GlobalRequestLoadBalancer 是一个全局共享的 Ray Actor,负责将 LLM 推理请求路由到具体的 vLLM Server Replica。所有 AgentLoopWorker 共享同一个 LoadBalancer 实例。

Global Load Balancer:粘性会话走 LRU,新请求走最少在途,Partial Rollout 重试仍落到同一 replica
Global Load Balancer 架构。共享 Ray Actor,sticky session + least inflight;Partial Rollout 用同一 request_id 复用 KV Cache。点击可打开原图。

2.5.1 核心数据结构

# agent_loop.py:65-97
@ray.remote
class GlobalRequestLoadBalancer:
    def __init__(self, server_actor_ids: list[str], max_cache_size: int = 10000):
        self._inflight_requests: dict[str, int] = {sid: 0 for sid in server_actor_ids}
        self._request_id_to_server: LRUCache = LRUCache(maxsize=max_cache_size)
        # 例如: _inflight_requests = {
        #   "replica_0": 12,   # replica_0 上当前有 12 个在途请求
        #   "replica_1": 10,   # replica_1 上当前有 10 个在途请求
        # }

2.5.2 acquire_server:路由决策算法

def acquire_server(self, request_id: str) -> str:
    # 策略 1: 粘性会话(Sticky Session)
    # 同一 request_id 的路由记录在 LRU Cache 中,保证同一对话的所有 turn 发送到同一 replica
    if request_id in self._request_id_to_server:
        server_id = self._request_id_to_server[request_id]
        self._inflight_requests[server_id] += 1  # 在途请求数 +1
        return server_id

    # 策略 2: 最少在途请求数(Least Inflight Requests)
    # 新请求选择当前负载最轻的 replica
    server_id = min(self._inflight_requests, key=self._inflight_requests.get)
    self._request_id_to_server[request_id] = server_id
    self._inflight_requests[server_id] += 1
    return server_id

def release_server(self, server_id: str) -> None:
    # 请求完成后在 finally 块中调用,在途请求数 -1
    self._inflight_requests[server_id] -= 1

2.5.3 在请求生命周期中的位置

AsyncLLMServerManager.generate() 中(被 FullyAsyncLLMServerManager.generate() 的每次内部迭代调用):

async def generate(self, request_id, *, prompt_ids, sampling_params, ...):
    # 1. 从 LoadBalancer 获取 server, 这里request_id是一个完整任务一个
    server_id, server = await self._acquire_server(request_id)
    #    → LoadBalancer.acquire_server(request_id)
    #      返回 (server_id, server_handle)

    try:
        # 2. 向选中的 vLLM Server 发送推理请求
        output = await server.generate.remote(
            request_id=uuid4().hex,  # 每次推理用新的 request_id,这个request_id是每轮对话不一样的
            prompt_ids=prompt_ids,
            sampling_params=sampling_params,
            ...
        )
        return output
    finally:
        # 3. 释放 server(fire-and-forget,减少延迟)
        self._release_server(server_id)
        #    → LoadBalancer.release_server(server_id)
        #      在途请求数 -1

2.5.4 设计目标

  1. 粘性会话(Sticky Session):同一对话的多个 turn 被分配到同一个 vLLM Server Replica,使得 vLLM 的 Automatic Prefix Caching(APC)可以复用前面 turn 的 KV Cache,显著减少 prefill 开销。
  2. 最少请求均衡(Least Inflight):新对话(首次请求)被分配到当前负载最轻的 replica,避免某个 replica 被过度分配而另一个闲置。

2.5.5 与 FullyAsyncLLMServerManager 的交互

FullyAsyncLLMServerManager.generate(request_id="task_abc")
  │
  ├── 第 1 次调用 super().generate(request_id="task_abc")
  │     → LoadBalancer.acquire_server("task_abc") → 首次路由 → replica_0
  │     → vLLM 推理,中途被 abort → stop_reason="aborted"
  │     → LoadBalancer.release_server("replica_0")  (fire-and-forget)
  │
  ├── [partial_rollout=True, 自动重试]
  │
  └── 第 2 次调用 super().generate(request_id="task_abc")
        → LoadBalancer.acquire_server("task_abc") → LRU Cache 命中 → replica_0
        → 同一个 replica,继续生成

关键点:FullyAsyncLLMServerManager 的 partial rollout 重试用的是同一个 request_id,所以 LoadBalancer 会保证它被分配到同一个 vLLM Server Replica,从而可以复用已有的 KV Cache。


3Harbor SWE Agent Loop 适配

3.1 BuiltinSWEAgentLoop 概述

BuiltinSWEAgentLoop 继承自 AgentLoopBase,是与 VeRL 框架的集成桥梁。其核心设计思想是:将完整的 Harbor Trial 执行过程封装为 VeRL 的一次 run() 调用

VeRL AgentLoopManager
  └── BuiltinSWEAgentLoop.run(sampling_params, **kwargs)
       ├── 从 kwargs 提取 task_path 和 extra_info
       ├── 获取/启动 _VLLMChatCompletionsProxy(单例)
       ├── 调用 _run_harbor_trial() → Harbor Trial.run()
       ├── 从 Proxy Session 提取轨迹 token_ids
       └── 构建并返回 AgentLoopOutput

3.2 与 ToolAgentLoop 的区别

特性 ToolAgentLoop BuiltinSWEAgentLoop
谁控制 Agent 循环 VeRL(状态机驱动) Harbor(Trial.run 内部)
工具执行位置 VeRL 进程内 Harbor 沙箱内(K8s Pod/Docker)
模型推理调用方式 直接调用 server_manager.generate() 通过 HTTP Proxy 间接触发
轨迹 token 来源 状态机自身累积 Proxy 会话捕获
支持 Partial Rollout 通过 FullyAsyncLLMServerManager 通过 Proxy → FullyAsyncLLMServerManager
验证方式 VeRL RewardModel Harbor Verifier(运行测试套件)

3.3 关键配置注入

BuiltinSWEAgentLoop 接收 harbor_cfg 字典,包含:

harbor_cfg = {
    "agent": {
        "name": "oh",                          # Agent 名称
        "import_path": "harbor_patch...:OpenHands",  # Agent 类导入路径
        "max_iterations": 100,                 # 最大迭代次数
        "tool_parser": "hermes",               # Tool Parser 类型
        "trials_dir": "./harbor_trials",       # 轨迹输出目录
        "max_retries": 2,                      # 最大重试次数
        "env": {                               # Agent 环境变量
            "LLM_BASE_URL": "由Proxy动态注入",
            "LLM_API_KEY": "...",
        }
    },
    "task": {"path": "/path/to/task"},
    "environment": {
        "import_path": "...:KubernetesEnvironment",
        "override_cpus": 1,
        "override_memory_mb": 8192,
    }
}

3.4 错误处理与重试策略

_run_harbor_trial() 中的重试逻辑:

for attempt in range(max_retries):
    ├── 创建新 session_id, 新 api_base, 新 trajectory_dir
    ├── 调用 Harbor Trial.run()
    ├── 根据异常类型决定行为:
    │   ├── AgentTimeoutError  → 不重试,取 reward 返回
    │   ├── ContextLengthExceededError → 不重试,reward=0
    │   ├── NonZeroAgentExitCodeError  → 不重试,取 reward 返回
    │   ├── 无 verifier_result  → 重试
    │   └── 其他异常           → 重试
    └── 成功 → break

4Harbor 内部 Rollout 执行过程

4.1 整体流程概览

一次完整的 Harbor Trial 执行包含四大阶段,对应四个 TimingInfo 对象:

┌──────────────────────────────────────────────────────────────────┐
│                    Harbor Trial.run()                             │
│                                                                   │
│  ┌─────────────────┐                                             │
│  │ 1. env_setup    │  K8s Pod 创建 / Docker 容器启动               │
│  │   环境搭建        │  执行任务 Dockerfile (RUN/WORKDIR/COPY)    │
│  └────────┬────────┘                                             │
│           ▼                                                       │
│  ┌─────────────────┐                                             │
│  │ 2. agent_setup  │  安装 Agent 依赖 / 配置环境变量               │
│  │   Agent 初始化   │  设置 LLM_BASE_URL 指向 Proxy              │
│  └────────┬────────┘                                             │
│           ▼                                                       │
│  ┌─────────────────────────────────────────────────────────────┐ │
│  │ 3. agent_execute (多轮 LLM 交互 + 工具执行循环)               │ │
│  │                                                              │ │
│  │   ┌──────────────────┐     ┌──────────────────┐             │ │
│  │   │  model_inference │────>│  tool_execute    │             │ │
│  │   │  HTTP POST 到    │<────│  在沙箱中执行      │             │ │
│  │   │  Proxy (/v1/chat │     │  bash/file/git等  │             │ │
│  │   │   /completions)  │     │                   │             │ │
│  │   └──────────────────┘     └──────────────────┘             │ │
│  │         重复直到 Agent 决定 stop 或达到限制                    │ │
│  └─────────────────────────────┬───────────────────────────────┘ │
│                                ▼                                  │
│  ┌─────────────────┐                                             │
│  │ 4. verify       │  运行测试套件 (pytest/unittest)              │
│  │   结果验证        │  计算 reward                               │
│  └─────────────────┘                                             │
└──────────────────────────────────────────────────────────────────┘

4.2 阶段详解

4.2.1 env_setup(环境搭建)

KubernetesEnvironment(或 RemoteDockerEnvironment)负责:

  1. 解析任务目录中的 Dockerfile
  2. 创建 K8s Pod(或 Docker 容器),应用资源配置(CPU/Memory/Volumes)
  3. 执行环境构建(Dockerfile 中的 RUN 命令)
  4. 支持 inline-build 模式:在 Pod 内直接执行 Dockerfile 指令
  5. 支持 Nydus 镜像加速和 hostPath 路径重写

4.2.2 agent_setup(Agent 初始化)

  1. 将 Agent 运行时挂载到环境(image volume 或 pip install)
  2. 设置环境变量(LLM_BASE_URL, LLM_API_KEY 等),指向 Proxy 的地址
  3. 渲染 Agent 的 user prompt(Jinja2 模板)

4.2.3 agent_execute(Agent 执行循环)

这是核心的多轮交互阶段。Harbor Agent(如 OpenHands)在沙箱中运行:

┌─────────────┐     HTTP POST /v1/chat/completions    ┌──────────────┐
│             │ ──────────────────────────────────>  │              │
│   Harbor    │     {"messages": [...], "tools": ..}  │   Proxy      │
│   Agent     │                                       │  (aiottp)    │
│  (LiteLLM)  │ <────────────────────────────────── │              │
│             │     {"choices": [{...}], ...}        │              │
└──────┬──────┘                                       └──────┬───────┘
       │                                                      │
       │ 收到 assistant message                               │ server_manager
       │ (含 content 和 可选的 tool_calls)                     │ .generate()
       ▼                                                      │
  ┌──────────┐                                               │
  │ 执行工具  │  bash, edit, git, etc.                         │
  └────┬─────┘                                               │
       │ 工具结果作为新 user/tool message                       │
       │ 加入消息历史                                          │
       └──────────────────────────────────────────────────────┘
                    下一轮 HTTP 请求

4.2.4 verify(结果验证)

  1. 在沙箱中运行测试套件(通过 Harbor Verifier)
  2. 返回二值 reward(0.0 或 1.0)

4.3 时序数据收集

*collect × harbor_timings()TrialResult 中提取各阶段耗时,上报到 VeRL 的 metrics 系统:

metrics = {
    "harbor_env_setup": 12.3,       # 环境搭建耗时 (秒)
    "harbor_agent_setup": 5.1,      # Agent 初始化耗时
    "harbor_agent_execute": 45.7,   # Agent 执行循环耗时
    "harbor_verifier": 30.2,        # 验证耗时
    "harbor_total": 93.3,           # 总耗时
}

5轨迹代理 Proxy 功能

5.1 整体架构

_VLLMChatCompletionsProxy 是一个进程内的 aiohttp HTTP Server,伪装成一个标准的 OpenAI 兼容 vLLM 端点,使得 Harbor Agent 无需任何修改即可通过 LiteLLM 与之通信。

aiohttp HTTP Server 是什么

aiohttp 是 Python 的异步 HTTP 库。HTTP Server 的意思是:在当前进程里用它真正监听一个端口、接收 HTTP 请求、返回 JSON——和 nginx / uvicorn 扮演的角色一样,只是服务嵌在 VeRL 的 AgentLoopWorker 进程内部,不另起进程。

Harbor 沙箱里的 Agent(LiteLLM / Claude Code)只会调「OpenAI 兼容的 HTTP 接口」,不会直接调 VeRL 的 Python 函数。所以 Worker 启动时用 aiohttp 绑在 0.0.0.0:{随机端口},对外宣称自己是 vLLM:POST /v1/chat/completions 进来后,Proxy 做 tokenize、记轨迹,再转给 server_manager.generate()

之所以用 aiohttp 而不是 Flask,是因为它和 VeRL 共用同一条 asyncio 事件循环:几十条并发轨迹的 HTTP 请求可以同时挂起等待推理,不会把事件循环堵死。

In-Process Proxy:Harbor Agent 经 aiohttp 假 vLLM 做 tokenize 与轨迹捕获,再转给 VeRL server_manager
In-Process Proxy 架构。进程内 aiohttp 伪装成 OpenAI 兼容 vLLM,做 token-in / token-out 与轨迹捕获。点击可打开原图。
┌─────────────────────────────── 进程边界 ───────────────────────────────┐
│                                                                        │
│  Harbor Agent                Proxy (aiohttp)            VeRL          │
│  (K8s Pod/Docker)          0.0.0.0:{ephemeral_port}   server_manager │
│       │                           │                          │        │
│       │  POST /sess/{id}/v1/      │                          │        │
│       │  chat/completions         │                          │        │
│       │ ────────────────────────> │                          │        │
│       │                           │  1. normalize messages   │        │
│       │                           │  2. apply_chat_template │        │
│       │                           │  3. compute prompt_ids  │        │
│       │                           │                           │        │
│       │                           │  server_manager.generate │        │
│       │                           │ ─────────────────────────>        │
│       │                           │                           │        │
│       │                           │  TokenOutput              │        │
│       │                           │ <─────────────────────────        │
│       │                           │                           │        │
│       │                           │  4. tool parse (hermes/   │        │
│       │                           │     qwen3_coder)          │        │
│       │                           │  5. update session state  │        │
│       │  OpenAI-format JSON       │                           │        │
│       │ <──────────────────────── │                           │        │
│       │                           │                           │        │
└────────────────────────────────────────────────────────────────────────┘

5.2 生命周期管理

Proxy 生命周期:

代理初始化 (per server_manager Ray Actor, 单例)
  ├── double-checked locking (threading.Lock)
  ├── 预计算 _gen_prompt_len (add_generation_prompt 的 token 数)
  ├── 预计算 _assistant_eos_tail (EOS 后的尾部 token)
  └── 创建 vLLM Tool Parser 实例

代理启动 (asyncio.Lock, 惰性启动)
  ├── 绑定 aiohttp 到 0.0.0.0:{OS 分配的端口}
  └── host IP (用于 api_base URL)

每次 Harbor Trial:
  ├── open_session(session_id, trajectory_dir)
  │   └── 初始化 session 状态 (traj_acc_ids, mask, logprobs 等)
  │
  ├── [Harbor 发送多个 HTTP 请求]
  │   └── _handle_chat() 处理每个请求
  │
  └── pop_session(session_id)
      ├── 读取累积的轨迹状态
      └── 写入 proxy_trajectory.json (用于调试)

5.3 Token-in-Token-out 保证

"Token-in-Token-out" 是 Proxy 最关键的保证——即增量式 tokenization 的结果必须与一次性 apply_chat_template 的结果逐 token 一致。这是因为 VeRL Trainer 需要精确知道每个 token 的定位。

5.3.1 正常增量路径

第 1 次请求:
  messages = [system, user]
  → apply_chat_template(all messages, tools)
  → traj_acc_ids = [token_ids]  (作为 prompt_ids)
  → initial_prompt_token_len = len(traj_acc_ids)
  → messages_snapshot = [system, user]
  → 生成完成后追加 assistant response (mask=1)

第 2 次请求:
  messages = [system, user, assistant, tool(user)]
  → messages 是 snapshot 的前缀扩展 ✓
  → delta = [tool(user)]  (只 tokenize 新增消息)
  → delta_ids = apply_chat_template(delta, remove_system_prompt=True)
  → traj_acc_ids += delta_ids  (mask=0, logprob=0.0)
  → 生成完成后追加 assistant response (mask=1)

第 3+ 次请求:
  → 同样只 tokenize delta,累积到 traj_acc_ids

5.3.2 EOS 尾部回填

vLLM 在生成 EOS token 后立即停止,不会输出 chat template 中 EOS 之后的 token(如 <|im_end|> 后面的换行符 \n)。Proxy 预计算了 assistant_eos_tail,在每次以 EOS 结束的生成后回填这些 token:

vLLM 输出: [token_1, token_2, ..., token_eos]
                                    ↑ 在此停止
chat template 完整输出: [token_1, token_2, ..., token_eos, \n]
                                                          ↑ 被省略

Proxy 回填: traj_acc_ids += [\n]  (mask=1, logprob=0.0)

5.3.3 前缀不匹配回退路径

当 Harbor Agent 修改了历史消息(不是简单追加),Proxy 会回退到逐段重建模式:

重建逻辑:
  prompt segment = apply_chat_template(prefix_msgs, tools)
  for each assistant: format(msg) 去掉头尾 gen_prompt (mask=1)
  for each non-assistant batch: apply_chat_template(batch) (mask=0)

  结果与增量路径逐 token 一致
  同时写入 debug JSON 文件用于排查
这段在解决什么

增量 tokenize 的前提坏了:不能再「只编新消息」,必须按增量路径同款切法把整段对话重新拼成一条 token 串,并且还能分出 mask=1(模型生成)和 mask=0(工具/用户观察)。

5.3.1 的增量路径在赌一件事:新请求的 messages 只是旧 messages_snapshot尾巴延长

第 1 次:  [system, user]
第 2 次:  [system, user, assistant, tool]          ← 前面原封不动,只在尾巴加
第 3 次:  [system, user, assistant, tool, assistant, tool]

只要这个成立,Proxy 就只 tokenize 新增的 delta,接到 traj_acc_ids 后面。

前缀不匹配 = 新来的 messages 不是旧 snapshot 的前缀扩展(_messages_is_prefix 失败),或 tools 变了(_tools_equal 失败)。常见触发:

这时「只编 delta」会编错位,整条轨迹对 Trainer 就废了。

不能偷懒改成一次 apply_chat_template(全部 messages)。一次编完能得到正确的 token 串,但分不出哪些 token 是模型生成的、哪些是工具观察。PPO 需要 assistant 内容 mask=1、system / user / tool mask=0。所以必须按轮切开再拼,一边出 token、一边打 mask。这就是「逐段重建」。

Chat template 里有一块固定的「请开始当助手说话」标记,叫 gen_prompt。Qwen 一般是 3 个 token:<|im_start|>assistant\napply_chat_template(...) 默认会在末尾再挂一块,等模型接着写。一段真实对话在 token 线上长这样:

[system][user] <|im_start|>assistant\n     ← prompt 段,以 gen_prompt 结尾
  你好吗<|im_end|>\n                         ← 模型生成(mask=1)
[tool 结果] <|im_start|>assistant\n          ← 观察段(mask=0),末尾再挂 gen_prompt
  我去改文件<|im_end|>\n                     ← 又一轮生成(mask=1)

重建按时间顺序走三步,刻意复刻增量路径的切法:

1. prompt segment = apply_chat_template(prefix_msgs, tools)

prefix_msgs第一条 assistant 之前的消息(通常是 system + 第一个 user,可能带 tools 定义)。编出来的串以一块 gen_prompt 结尾。这就是 initial_prompt,对应 Trainer 里的 prompt,不算 response。

2. 每条 assistant:单独编,去掉头尾 gen_prompt(mask=1)

单独把一条 assistant 丢进 apply_chat_template,会得到:

asst_ids == <|im_start|>assistant\n + 正文 + <|im_end|>\n + <|im_start|>assistant\n
               ↑ 头:上一段末尾已经有了                    ↑ 尾:会和下一段开头撞车

头已经在上一段末尾,尾会变成「下一轮 assistant 的开头」重复一次。所以源码剥掉两端各 gp_len 个 token:

asst_ids = asst_ids[gp_len:-gp_len]

剥完只留下正文 + EOS 尾巴,打 mask=1

3. 连续的非 assistant:整批一起编(mask=0)

中间的 user / tool 不要一条条编。增量路径本来就是 delta = 连续新增的非 assistant 一次编完。如果重建时拆开编,模板可能多插分隔符,和增量路径对不齐。所以遇到连续的 tool/user,整批一起编,得到「观察内容 + 末尾 gen_prompt」,打 mask=0。末尾那块 gen_prompt 正好当下一轮 assistant 的开头。

走完之后,token 序列和「从来没改过历史、一直走增量路径」是逐 token 相同的(再补上 5.3.2 那个 EOS 后的 \n)。

增量路径(历史只追加) 回退重建(历史被改过)
第 1 次:整段 prompt 一次编 prefix_msgs 一次编
生成后把 assistant token 接上 每条 assistant 单独编,剥头尾 gen_prompt
以后只编 messages[len(snap):] 连续非 assistant 按批编

回退是异常路径。Proxy 会在 trajectory_dir 写下三份 debug JSON,用来 diff 到底是哪一条被改、删、换了:

重建后 logprob 怎么处理

重建后 logprob 先填 0。如果其实是 abort/retry——新历史只是旧轨迹的前缀——公共前缀上的 logprob 仍然有效,会从旧 archive 拷回来;分叉点之后继续为 0。所以重建不等于「整条轨迹 logprob 全丢」,多数 retry 场景只丢分叉之后那截。

5.4 Rollout LogP 和 Expert Choice 传回

在 *finalize × generation() 中:

# 从 vLLM TokenOutput 提取 log_probs
lp = output.log_probs
if lp is None:
    lp = [0.0] * len(output_ids)

# 累积到 session 状态
s["traj_response_mask"] += [1] * len(output_ids)     # LLM 生成的 token
s["traj_response_logprobs"] += list(lp)              # 每个 token 的 logp

# routed_experts 在 FullyAsyncLLMServerManager 中被累积
# 用于 MoE 路由重放
if output.routed_experts is not None:
    final_output.routed_experts = torch.cat([...])

最终这些数据通过 AgentLoopOutput 传递给 Trainer:

5.6 独占模式下的早停机制

Proxy 支持 max_consecutive_no_tool 参数:当连续 N 个 assistant turn 都没有 tool_call 时,强制终止 session。这用于防止 Agent 进入"对话死循环"。


6VeRL 训练侧流程

6.1 FullyAsyncTrainer 整体流程

FullyAsyncTrainer 是训练的消费者,继承自 SeparateRayPPOTrainerRayPPOTrainer

┌─────────────────────────────────────────────────────────────┐
│                 FullyAsyncTrainer.fit()                      │
│                                                              │
│  while True:                                                 │
│    ┌──────────────────────────────────────────────────┐    │
│    │              fit_step()                           │    │
│    │                                                   │    │
│    │  1. _fit_generate                                │    │
│    │     └─ _get_samples_from_queue()                │    │
│    │        从 MessageQueue 逐个取样本                  │    │
│    │        攒够 required_samples 后 assemble batch    │    │
│    │                                                   │    │
│    │  2. _fit_compute_reward                          │    │
│    │     └─ reward_model 评分 (Harbor 已完成 reward)  │    │
│    │                                                   │    │
│    │  3. _fit_compute_log_prob                        │    │
│    │     └─ _compute_old_log_prob()                  │    │
│    │        重算 old_log_prob (含 MIS 支持)           │    │
│    │                                                   │    │
│    │  4. _fit_compute_ref_log_prob                    │    │
│    │     └─ Reference Policy 的 log prob              │    │
│    │                                                   │    │
│    │  5. _fit_compute_critic                          │    │
│    │     └─ Critic 模型估值                            │    │
│    │                                                   │    │
│    │  6. _fit_compute_advantage                       │    │
│    │     └─ GAE 优势函数计算                           │    │
│    │                                                   │    │
│    │  7. _fit_update_critic                           │    │
│    │     └─ Critic 参数更新                            │    │
│    │                                                   │    │
│    │  8. _fit_update_actor                            │    │
│    │     └─ Actor PPO/GRPO/DAPO loss 更新             │    │
│    │                                                   │    │
│    │  9. _fit_update_local_step                       │    │
│    │     └─ local_trigger_step 递增                    │    │
│    │        达到 trigger_parameter_sync_step 则复位     │    │
│    │                                                   │    │
│    │  10. _fit_update_weights                         │    │
│    │     └─ 若 local_trigger_step==1:                │    │
│    │        → checkpoint_manager.update_weights()    │    │
│    │        → rollouter.reset_staleness()            │    │
│    └──────────────────────────────────────────────────┘    │
│    if TrainingStopException: break                          │
└─────────────────────────────────────────────────────────────┘

6.2 样本获取流程

async def _get_samples_from_queue(self):
    """从 MessageQueue 逐个取样本,攒够 required_samples 个"""
    samples = []
    for _ in range(self.required_samples):
        raw_sample = await self.message_queue_client.get_sample()
        if raw_sample is None:  # 终止信号
            break
        rollout_sample = ray.cloudpickle.loads(raw_sample)
        samples.append(rollout_sample)

    # 组装成 DataProto batch
    batch = assemble_batch_from_rollout_samples(samples)
    return batch

6.4 Rollout Importance Sampling (IS) 与 bypass_mode

bypass_mode = True (默认):
  old_log_probs = rollout_log_probs  (零额外计算)

bypass_mode = False:
  old_log_probs = 用训练引擎 (FSDP/Megatron) 重新计算
  同时启用 Rollout IS 修正 (compute_rollout_correction_weights)
  等价于 AReaL 的 Decoupled PPO

6.5 local_trigger_step 与参数同步时机

local_trigger_step 计数器逻辑:

fit_step 1  → local_trigger_step = 1 → 触发 weight sync + reset_staleness
fit_step 2  → local_trigger_step = 2
fit_step 3  → local_trigger_step = 3
...
fit_step N  → local_trigger_step = trigger_parameter_sync_step 
            → 下次复位为 1 并 current_param_version += 1

这种设计实现了"多步本地训练 + 一次参数同步"的 off-policy pipeline。


7权重同步整体流程与 Partial Rollout

7.1 权重同步完整流程

CheckpointEngineManager.update_weights() 是权重同步的核心入口:

┌─────────────────────────────────────────────────────────────────┐
│          CheckpointEngineManager.update_weights()                │
│                                                                  │
│  Step 1: Abort all in-flight requests                           │
│    └─ 对每个 Rollout Replica 调用 replica.abort_all_requests()  │
│       → vLLM 收到 abort 信号,当前推理终止                        │
│       → TokenOutput.stop_reason = "aborted"                     │
│                                                                  │
│  Step 2: 构建临时 RayWorkerGroup                                │
│    └─ 收集所有 Rollout Replica 的 workers                        │
│                                                                  │
│  Step 3: Sleep replicas (可选,free_cache_engine)               │
│    └─ 释放 vLLM KV-cache GPU 显存                                │
│       → 为 NCCL 权重传输释放更多的 GPU 显存空间                   │
│                                                                  │
│  Step 4: 构建 NCCL Process Group                               │
│    └─ Trainer rank 0 + 所有 Rollout workers                     │
│       → 跨进程的 NCCL 通信域                                     │
│                                                                  │
│  Step 5: 权重传输                                               │
│    ├─ Trainer rank 0: send_weights()                           │
│    │   ├─ 遍历模型 weight generator                            │
│    │   ├─ 将参数 pack 到 CuPy send buffer (双缓冲)              │
│    │   └─ 通过 NCCL Broadcast 发送                              │
│    │                                                             │
│    ├─ Rollout workers: receive_weights()                       │
│    │   ├─ 通过 NCCL Broadcast 接收                              │
│    │   ├─ 接收 ZMQ PUB/SUB 元数据 (tensor name, shape, offset) │
│    │   └─ Unpack 到目标 tensor                                  │
│    │                                                             │
│    └─ ZeroMQ PUB/SUB 与 NCCL 并行传输元数据                     │
│                                                                  │
│  Step 6: Finalize                                               │
│    └─ 释放 NCCL bucket, 销毁 process group (可选)                │
│                                                                  │
│  Step 7: Wake up replicas                                       │
│    └─ 恢复 vLLM KV-cache, 重新加载权重到 GPU                     │
│                                                                  │
│  Step 8: Resume unfinished requests (Partial Rollout)          │
│    └─ 对每个 Replica 调用 replica.resume_generation()          │
│       → 中断的样本从中断点继续生成                                │
└─────────────────────────────────────────────────────────────────┘

7.2 Partial Rollout 机制详解

无 Partial · 等全部跑完
|-- 样本1 --|-- 样本2 --|-- 样本3 --|-- 等待 --|-- 同步 --|-- 样本4 --|
                                      ↑
                              同步前必须清空 active_tasks
有 Partial · 切断再续写
|-- 样本1 --|-- 样本2 --|-- 样本3…中断…同步后继续…|-- 样本4 --|
                          ↑                    ↑
                   abort_all_requests()   resume_generation()

Partial Rollout 解决了权重同步时的一个关键问题:如何在不停掉进行中的样本生成的情况下完成权重同步?

7.2.1 无 Partial Rollout(staleness_threshold > 0, partial_rollout = False)

同步到来时先等待所有 active_tasks 结束,再传权重。吞吐被长尾样本拖住。

7.2.2 有 Partial Rollout(staleness_threshold > 0, partial_rollout = True)

abort_all_requests() 让当前推理以 stop_reason=aborted 结束;同步完成后 resume_generation() 从中断点继续,对 Agent Loop 仍是一次 generate()

7.2.3 FullyAsyncLLMServerManager 的中断恢复实现

class FullyAsyncLLMServerManager:
    async def generate(self, ...):
        final_output = TokenOutput(token_ids=[], ...)

        while True:
            # 每次迭代传入 prompt_ids + 已生成的 tokens
            output = await super().generate(
                prompt_ids=prompt_ids + final_output.token_ids,
                ...
            )

            # 累积到 final_output
            final_output.token_ids.extend(output.token_ids)
            final_output.log_probs.extend(output.log_probs)

            # 检查是否需要继续(中断恢复)
            if output.stop_reason in ("aborted", "abort") \
               and partial_rollout is True:
                continue  # 继续循环,从中断点重新生成
            else:
                break  # 正常完成,退出循环

        return final_output

关键点:

7.3 权重同步时序图

Trainer                        ParameterSynchronizer           Rollouter/Replicas
  │                                    │                            │
  │ _fit_update_weights()              │                            │
  │  (local_trigger_step == 1)         │                            │
  │───────────────────────────────────>│                            │
  │                                    │ abort_all_requests()       │
  │                                    │───────────────────────────>│
  │                                    │                            │ 正在生成的请求
  │                                    │                            │ 收到 abort 信号
  │                                    │                            │ stop_reason="aborted"
  │                                    │                            │
  │                                    │ sleep_replicas()           │
  │                                    │───────────────────────────>│ 释放 KV-cache
  │                                    │                            │
  │                                    │ build NCCL process group   │
  │                                    │<─ Trainer + Replicas ─────>│
  │                                    │                            │
  │ send_weights()                     │                            │
  │ (NCCL Broadcast)                   │                            │
  │════════════════════════════════════════════════════════════════>│ receive_weights()
  │                                    │                            │
  │                                    │ wake_up_replicas()         │
  │                                    │───────────────────────────>│ 恢复 KV-cache
  │                                    │                            │
  │                                    │ resume_generation()        │
  │                                    │───────────────────────────>│ 继续未完成的请求
  │                                    │                            │
  │ rollouter.reset_staleness()        │                            │
  │────────────────────────────────────────────────────────────────>│
  │                                    │                            │ paused=False
  │                                    │                            │ staleness_samples=len(active_tasks)

7.5 Reset Staleness 逻辑

同步完成后,Trainer 调用 rollouter.reset_staleness()

async def reset_staleness(self):
    async with self.lock:
        # 将在途任务数 + 队列中等待的样本数 设为新的 staleness 基准
        self.staleness_samples = len(self.active_tasks) + self.pending_queue.qsize()

        # 解除暂停信号
        self.paused = False
        self._resume_event.set()

        # 记录 idle/active 时间用于计算 idle_ratio

这样设计的原因:staleness_samples 需要反映"在当前参数版本之后生成的样本数"。reset 后,已有的 active_tasks 和 pending 样本都算作新版本下的样本。


A完整链路流程总图

以下是 fully_async_2nodes.sh 启动后的完整数据流:

                           ┌───────────────────┐
                           │  fully_async_main   │  (Hydra 入口)
                           │   (FullyAsyncTaskRunner) │
                           └─────────┬─────────┘
                                     │
                     创建 MessageQueue, FullyAsyncTrainer, FullyAsyncRollouter
                     初始参数同步, Trainer ← NCCL → Rollouter
                                     │
              ┌──────────────────────┼──────────────────────┐
              ▼                      ▼                      ▼
    ┌─────────────────┐   ┌──────────────┐   ┌─────────────────────┐
    │ FullyAsyncTrainer │   │ MessageQueue  │   │ FullyAsyncRollouter │
    │   (1 node, 8 GPU)  │   │  (Ray Actor)  │   │   (1 node, 8 GPU)   │
    └────────┬────────┘   └──────┬───────┘   └──────────┬──────────┘
             │                   │                       │
             │   get_sample()    │      put_sample()     │
             │<──────────────────│<──────────────────────│
             │                   │                       │
    ┌────────▼────────┐          │              ┌────────▼──────────┐
    │ PPO Training     │          │              │ _feed_samples()   │
    │ 1. get samples   │          │              │  从 dataloader     │
    │ 2. reward        │          │              │  逐条读取任务       │
    │ 3. old_log_prob  │          │              │  → pending_queue   │
    │    (MIS支持)     │          │              └────────┬──────────┘
    │ 4. ref_log_prob  │          │                       │
    │ 5. critic        │          │              ┌────────▼──────────┐
    │ 6. advantage     │          │              │ _processor_worker()│
    │ 7. update critic │          │              │  逐个处理样本       │
    │ 8. update actor  │          │              │  max_concurrent     │
    │                  │          │              │  控制并发数         │
    │ sync_weights()──>│          │              └────────┬──────────┘
    │                  │          │                       │
    │                  │          │        ┌──────────────▼──────────────┐
    │                  │          │        │ _process_single_sample()    │
    │                  │          │        │   → async_rollout_manager   │
    │                  │          │        │     .generate_sequences_    │
    │                  │          │        │     single()                │
    │                  │          │        └──────────────┬──────────────┘
    │                  │          │                       │
    │                  │          │        ┌──────────────▼──────────────┐
    │                  │          │        │ FullyAsyncAgentLoopManager  │
    │                  │          │        │   → Worker Queue            │
    │                  │          │        │   → BuiltinSWEAgentLoop     │
    │                  │          │        │     .run()                  │
    │                  │          │        └──────────────┬──────────────┘
    │                  │          │                       │
    │                  │          │        ┌──────────────▼──────────────┐
    │                  │          │        │ Harbor Trial.run()          │
    │                  │          │        │  ├─ env_setup (K8s Pod)     │
    │                  │          │        │  ├─ agent_setup             │
    │                  │          │        │  ├─ agent_execute loop:    │
    │                  │          │        │  │   ├─ HTTP → Proxy        │
    │                  │          │        │  │   │  → _handle_chat()   │
    │                  │          │        │  │   │  → server_manager   │
    │                  │          │        │  │   │    .generate()      │
    │                  │          │        │  │   │  → (FullyAsync)     │
    │                  │          │        │  │   └─ tool_execute       │
    │                  │          │        │  └─ verify                 │
    │                  │          │        └─────────────────────────────┘
    │                  │          │
    │<── AgentLoopOutput ─────────┘  (经 MessageQueue 传递)
    │    prompt_ids, response_ids,
    │    response_mask, response_logprobs,
    │    reward_score, num_turns,
    │    min_global_steps, max_global_steps

B关键文件索引

文件 作用
fully_async.md 全异步训练官方文档
fully_async_2nodes.sh 2 节点全异步启动脚本
fully_async_main.py 全异步入口,FullyAsyncTaskRunner 编排
fully_async_trainer.py FullyAsyncTrainer:训练消费端
fully_async_rollouter.py FullyAsyncRollouter:样本生产端
message_queue.py MessageQueue + MessageQueueClient
agent_loop.py AgentLoopBase + AgentLoopWorker + AgentLoopManager
tool_agent_loop.py ToolAgentLoop 状态机实现
agent_loop.py (async) FullyAsyncLLMServerManager + FullyAsyncAgentLoopManager
base.py CheckpointEngineManager:权重同步编排
nccl_checkpoint_engine.py NCCLCheckpointEngine:NCCL 权重传输
builtin_swe_agent_loop.py Harbor SWE Agent Loop 适配
vllm_chat_completion_proxy.py 进程内 vLLM HTTP Proxy
harbor_verl_fully_async_fsdp.yaml 全异步 FSDP 训练配置

C术语表

术语 解释
Colocate 训推一体架构,Rollout 和 Training 共享同一组 GPU
Fully Async 全异步训推分离架构,Rollouter 和 Trainer 完全解耦
Rollouter 样本生成器,持续异步地产生训练样本
MessageQueue 基于 Ray Actor 的异步消息队列,连接 Rollouter 和 Trainer
Staleness 陈旧度,指样本生成时所用参数与当前训练参数之间的版本差异
Partial Rollout 中断恢复机制,权重同步时暂停进行中的 rollout,同步后从中断点继续
MIS Multiple Importance Sampling,多步重要性采样,保证 pi_old 一致性
bypass_mode 直接使用 rollout_log_probs 作为 old_log_probs,省去重算
Token-in-Token-out 增量式 tokenization 与一次性 tokenization 结果逐 token 一致
Router Replay MoE 路由器决策的复现能力,需要记录 routed_experts 和 global_steps
Proxy 进程内 HTTP Server,将 Harbor 的 OpenAI 格式请求转发到 VeRL 的 server_manager
AgentLoop VeRL 中封装单条轨迹生成逻辑的抽象单元
checkpoint-engine 基于 NCCL 的高效分布式参数同步引擎
回到顶部 ← 全部解读