主入口 main_ppo.py 已迁到 V1:TransferQueue 数据仓库 + AgentLoop 服务化 rollout + 可插拔 trainer 模式,
默认 trainer.use_v1=true。本文只回答三个你真正关心的问题——为什么改、老脚本还能不能跑、我的自由度提高了多少。
use_v1=true + trainer_mode=sync ≈ 老同步 PPO;常规 GRPO/PPO 脚本基本原样可跑。想要异步只是加几行配置。一句话:老 pipeline 的三条隐含假设(同步、定长、单轮)挡住了 Agent 训练。大改就是把这三条拆掉。
RayPPOTrainer。{uid}_{sid}_{idx}(uid=原始 prompt,sid=第几条 session,idx=同一 loop 第几个 output),tag 里带 staleness 元数据。PPOTrainer.step() 与 TransferQueue 之间是双向箭头(读字段 → 计算 → 写回),各阶段都在仓库里进出。换句话说,V1 把老的「大 trainer」拆成了三个各司其职的组件:
生成侧(AgentLoopManager + LLM Server)、
数据桥(TransferQueue)、
训练侧(PPOTrainer + CheckpointEngine)。这也解释了为什么后面两问的答案都很干净——组件解耦了,你能单独替换其中任意一块。
python3 -m verl.trainer.main_ppo ...。
默认 trainer.use_v1=true + trainer.v1.trainer_mode=sync,行为≈老的同步 PPO。
仓库里 examples/grpo_trainer/*.sh 这些脚本压根没写 use_v1 / trainer_mode,照样跑。
要的话临时退回旧实现:trainer.use_v1=false(走 main_ppo_v0.py),但它 deprecated、计划 v0.9.0 删除,别长期依赖。
| 配置项 | 默认 | 作用 |
|---|---|---|
| trainer.use_v1 | true | 走 V1 主线;false 回退旧 main_ppo_v0.py(将删) |
| trainer.v1.trainer_mode | sync | sync / colocate_async / separate_async |
| trainer.v1.sampler.max_off_policy_threshold | 8 | 一条轨迹最多允许跨几个权重版本(异步时的陈旧度闸门) |
| trainer.v1.sampler.max_off_policy_strategy | drop | 超阈值怎么办:drop 丢弃 / wait 等待 |
| trainer.v1.colocate_async.num_warmup_batches | 1 | 异步前预热多少批生成 |
| trainer.v1.separate_async.num_warmup_batches | 4 | 分离异步的预热批数 |
| trainer.v1.separate_async.parameter_sync_step | 4 | 每几步把新权重同步给 rollout |
| actor_rollout_ref.rollout.agent.default_agent_loop | single_turn_agent | 默认用哪种 agent loop(单轮/工具/自定义) |
同一个脚本:同步照跑;想开异步只加两行
# 常规同步(等价旧行为,无需任何改动) python3 -m verl.trainer.main_ppo \ algorithm.adv_estimator=grpo \ data.train_files=... actor_rollout_ref.model.path=... \ trainer.n_gpus_per_node=8 trainer.nnodes=1 # ↑ 默认已是 trainer.use_v1=true, trainer.v1.trainer_mode=sync # 想开异步(共卡 + partial rollout):追加两行即可 trainer.v1.trainer_mode=colocate_async \ trainer.v1.colocate_async.num_warmup_batches=1
verl/interactions、curriculum sampler、dynamic dataset、旧 tool examples:这些在 #6067(worker→engine)、#6302 等 [BREAKING] 里被废弃/迁移,需要改到 engine 与 AgentLoop 体系。AgentLoop(见 Q3 骨架)。RewardLoopManager(规则或模型),配置路径在 reward.reward_model.*;确认你的老 reward 挂载方式仍匹配。separate_async(分卡全异步)有硬约束:train_batch_size == ppo_mini_batch_size、checkpoint_engine.backend 必须非 naive(nccl/nixl/mooncake)、colocate reward 不支持(需 reward.reward_model.enable_resource_pool=true)。| trainer_mode | train / rollout 布局 | partial rollout | 权重同步时机 | 选它当你… |
|---|---|---|---|---|
| sync | colocate(共享显存) | 关闭 | 每步:sample 后 sleep,train 后 update | 单轮/短轨迹、要稳要复现(默认,零改动) |
| colocate_async | colocate | 开启 | 每步 update + resume_generation;带 warmup |
共卡想提吞吐、容忍轻度 off-policy |
| separate_async | 物理分卡 + 空闲切换 | 开启 | 按 parameter_sync_step 周期同步 |
大规模 / 长 agent 轨迹、要最高吞吐 |
parameter_sync_step 的权重同步。吞吐最高,但一个 batch 里可能混着多个权重版本的样本,需要 decoupled-PPO 修正。on_sample_end / on_step_end 两个钩子里怎么调 CheckpointEngine。一个容易被忽略、但直接影响你申请多少卡的点:三种模式里真正改变 GPU 分配的只有 separate_async。sync 和 colocate_async 用的是同一批卡、同样数量,差别只在调度时机(上一节的钩子)。
rollout.* 配置决定大小的独立 rollout 池,两组卡真正同时跑。| 模式 | GPU 池 | rollout 放哪 | 决定卡数的配置 |
|---|---|---|---|
| sync / colocate_async | 1 个(共卡) | 训练卡上 hybrid,时分复用 | trainer.n_gpus_per_node × nnodes |
| separate_async | 2 个(分卡) | 独立 rollout 池 + 训练卡空闲兼职 | 训练池同上 + rollout.n_gpus_per_node × rollout.nnodes |
这三条不是随意限制,全部由「rollout 跑在另一组卡上、且永不停机」这一点直接推导出来。
| 硬约束 | 根因:为什么必须这样 |
|---|---|
train_batch_size == ppo_mini_batch_size |
全异步要求每个训练步恰好推进一个策略版本——一次 sample() → 一次梯度更新 → 同步一次权重。这样 global_steps 才是忠实的「策略版本时钟」,陈旧度 (global_steps − 生成时的 step) / parameter_sync_step 才算得准。若一批被拆成多个 mini-batch 做多次 SGD,一个 step 内策略就跳了好几个版本,异步 staleness 记账失真、解耦 PPO 的「行为策略 ↔ 近端策略」对应关系也乱了。所以强制一批 = 一个 mini-batch = 一次更新。 |
checkpoint_engine.backend != naive |
naive(即 ColocatedCheckpointEngine)只做「同卡 colocated 的 all_gather」,依赖 actor 与 rollout 在同一批卡 / 同一 torch.distributed 通信域——这正是共卡模式的玩法。分卡后 rollout 在另一组卡、还不停机,必须换成能跨设备 / 跨节点把权重传过去的后端:nccl/hccl(all_gather + broadcast)、nixl/mooncake(p2p),它们本就是为「actor/rollout 分离集群」设计的。 |
colocate reward 不支持(需 enable_resource_pool=true) |
colocate reward 是靠「rollout 睡觉、腾出显存」这个窗口来临时加载 reward model 打分的。separate_async 的独立 rollout 卡从不暂停释放显存,压根没有这个窗口,所以 reward model 必须自己占一组独立卡(standalone reward 池)。而共卡两种模式能让 rollout sleep 腾显存,才允许 colocate reward。 |
V1 主循环本质是生产者 / 消费者流水线:每个 step 里 _add_batch_to_generate()(生产者,fire-and-forget,非阻塞)塞入 1 批 prompt;replay_buffer.sample()(消费者,阻塞轮询)等到「完成的轨迹」够 1 批才取走(还优先取最旧的以压 staleness)。生产与消费解耦——这步塞进去的,不一定是这步消费的那批。
num_warmup_batches 批,把带子填满;此后每步「塞 1、取 1」维持这个提前量,训练一结束就有完成的轨迹可取,rollout 也不空转。sample() 取走够一批的完成轨迹(靠提前量基本不空等)。on_sample_end:abort_replicas() 中止未完成请求(partial rollout 把半截轨迹存下来)+ sleep_replicas()(rollout 交出显存/KV cache)。on_step_end:update_weights()(推新权重)+ resume_generation_replicas()(唤醒 rollout,带新权重接着跑那些被 abort 的半截轨迹)。trajectory_spans > 1、跨权重版本)+ 提前量让队列不空。代价是这步样本多由上一版权重生成 → 轻度 off-policy。| 模式 | num_warmup_batches 默认 | 为什么是这个深度 |
|---|---|---|
| sync | 无(0) | 刻意不重叠:生成完这批立刻训练这批,最 on-policy |
| colocate_async | 1 | 共卡只能浅浅重叠一步,约 1 步 off-policy |
| separate_async | 4(+ parameter_sync_step=4) | 分卡全程并行,流水线更深、容忍更强 staleness |
一句话:预热深度 ≈ 流水线深度 ≈ off-policy staleness 预算。预热建立提前量、partial rollout 让长轨迹跨步续跑,两者合起来才让共卡的 colocate_async 在同一批卡上把生成与训练错峰重叠。
sync 跑通(几乎零改动);发现生成把卡闲置严重、且能接受轻度 off-policy,再上 colocate_async;只有当轨迹很长 / 规模很大、且愿意配非 naive 权重后端时,才上 separate_async。
V1 把「生成逻辑 / 采样策略 / 调度模式 / PPO 每个阶段」都做成了可插拔扩展点。 绝大多数定制,你只需写一个子类 + 一行配置,不用碰主循环。
| 你想改的东西 | 怎么接入 | 入口 / 配置 |
|---|---|---|
| 生成 / 多轮 / 工具逻辑 | 写一个 AgentLoop 子类 | @register + rollout.agent.agent_loop_config_path |
| 整个 rollout 调度器 | 自定义 AgentLoopManager | rollout.agent.agent_loop_manager_class |
| 采样 / off-policy 策略 | 自定义 ReplayBuffer 子类 | trainer.v1.sampler.custom_sampler.{path,name} |
| 调度模式(同步/异步布局) | 注册新 trainer_mode | @register_trainer + 生命周期钩子 |
| PPO 单个阶段(adv / logprob / …) | override PPOTrainer._compute_* | 子类 + register_trainer |
| off-policy 修正(IS / decoupled) | 配置 rollout_correction | algorithm.rollout_correction.* |
| reward(规则 / 模型) | RewardLoopManager / reward server | reward.reward_model.* |
PPO 每个阶段都是独立方法——step() 里按顺序调用,想改哪步就在 trainer 子类里覆盖对应方法,其它步骤不受影响:
# verl/trainer/ppo/v1/trainer_base.py :: PPOTrainer.step() _add_batch_to_generate → replay_buffer.sample → [_compute_reward_colocate] → _balance_batch → _compute_old_log_prob → [_compute_ref_log_prob] → [_compute_values] → _compute_advantage → [_update_critic] → _update_actor # 生命周期钩子(三种 trainer_mode 就靠它们区分): on_train_begin / on_step_begin / on_sample_begin / on_sample_end / on_step_end
下面三张卡片是最常见的三类定制,点开看最小代码骨架(基于当前源码 API,实验阶段接口可能微调)。
内置参考实现:single_turn_agent_loop.py(单轮)、tool_agent_loop.py(ReAct 多轮 + 工具)。你只需实现 run():
# my_agent_loop.py from typing import Any from verl.experimental.agent_loop.agent_loop import ( AgentLoopBase, AgentLoopOutput, register, ) @register("my_agent") # 注册名,供 config / 数据集按名选用 class MyAgentLoop(AgentLoopBase): async def run(self, sampling_params: dict[str, Any], **kwargs) -> AgentLoopOutput: messages = list(kwargs["raw_prompt"]) prompt_ids = await self.apply_chat_template(messages) # 多轮就在这里 while 循环:生成 → 解析工具调用 → 追加观测 → 再生成 out = await self.server_manager.generate( # 向 LLM Server 发请求 request_id=uuid, prompt_ids=prompt_ids, sampling_params=sampling_params, ) return AgentLoopOutput( prompt_ids=prompt_ids, response_ids=out.token_ids, response_mask=[1] * len(out.token_ids), # 工具返回段可置 0,不参与 loss num_turns=2, )
接入(Hydra 注入,无需改主代码):
# my_agents.yaml
- name: my_agent
_target_: my_agent_loop.MyAgentLoop
actor_rollout_ref.rollout.agent.agent_loop_config_path=my_agents.yaml \ actor_rollout_ref.rollout.agent.default_agent_loop=my_agent
# my_sampler.py from verl.trainer.ppo.v1.replay_buffer import ReplayBuffer class MyReplayBuffer(ReplayBuffer): def sample(self, global_steps: int, partition_id: str, batch_size: int): # 自定义优先级 / 难度课程 / staleness 筛选 batch, metrics = super().sample(global_steps, partition_id, batch_size) return batch, metrics
接入:
trainer.v1.sampler.custom_sampler.path=my_sampler.py \
trainer.v1.sampler.custom_sampler.name=MyReplayBuffer \
trainer.v1.sampler.max_off_policy_threshold=8 # 内建陈旧度闸门仍可叠加
# my_trainer.py from verl.trainer.ppo.v1 import PPOTrainer, register_trainer @register_trainer("my_mode") class MyTrainer(PPOTrainer): def on_sample_end(self): self.checkpoint_manager.sleep_replicas() # 采样完让出显存 def on_step_end(self): if self.global_steps % 2 == 0: # 例:每 2 步才同步一次权重 self.checkpoint_manager.update_weights(self.global_steps)
接入:trainer.v1.trainer_mode=my_mode(确保该模块被 import 以触发注册)。想改某个 PPO 阶段,同理在这里 override _compute_advantage 等方法即可。
如果你也看过 slime(清华 / 智谱)或 AReaL(蚂蚁 / 清华 IIIS),会发现三者在解同一道题—— 解耦生成与训练、加一个数据桥、异步同步权重、控制 staleness。差异主要在抽象哲学和解耦彻底程度。
TransferQueue==0.1.8 是可复用的独立组件(Ascend 开源);
verl 的异步 off-policy 修正直接指向 AReaL 论文——trainer_separate_async.py 里写着
# TODO: Support Decoupled PPO: https://arxiv.org/abs/2505.24298。三者设计是收敛的。
| 维度 | verl V1 | slime | AReaL |
|---|---|---|---|
| 数据桥 | TransferQueue | Data Buffer | Replay Buffer(用一次) |
| 解耦 / 布局 | sync / colocate_async / separate_async 三档 | 一个 --colocate 开关 |
完全分离集群(fully async 优先) |
| off-policy 控制 | max_off_policy_threshold + rollout_correction |
over-sample + 部分轨迹回收(APRIL) | max_head_offpolicyness + decoupled PPO |
| trainer 抽象哲学 | 结构化基类 + 注册表 + 钩子 | 刻意不封装,循环裸露在 train.py |
algorithm-first(AReaL-lite) |
抽象哲学:极简。slime 明确不用 trainer 类包裹,训练循环直接暴露在 train.py;靠移动 ray.get 的位置切同步/异步,靠一个 --colocate 决定共卡/分卡。与 verl「结构化基类 + 注册表 + 钩子」正好相反——一个偏灵活裸写,一个偏工程可插拔。
长尾优化(APRIL):over-sample(要 32 发 64)→ 够了 abort 其余 → 半截轨迹存 buffer、下轮续跑。对应 verl 的 partial rollout + min/max_global_steps 跨版本追踪。
来源: THUDM/slime ↗ · 介绍博客 ↗
fully asynchronous。rollout worker 不等待、持续流式;trainer 攒够一批就更新,更新后同步权重回 rollout。一个 batch 可能混着不同模型版本的样本。论文报告相比同步系统最高 2.77× 加速且性能不降。
两把稳定器:rollout.max_head_offpolicyness(=0 退化为同步,异步常用 2–8)+ actor.use_decoupled_loss(解耦行为策略 π_behav 与近端策略 π_prox)。
与 verl 的血缘:verl separate_async 的 TODO 直指这篇;verl 的 rollout_correction(token-IS / seq-IS / geo-RS 等 decoupled_* 预设)与 bypass_mode,本质就是这套修正在 verl 的落地。
trainer.v1.trainer_mode=colocate_async(或规模够大再上 separate_async + 非 naive 权重后端)。#6067 迁到 engine;use_v1=false 只是临时退路(v0.9.0 移除)。