Harbor × VeRL · 全异步强化学习

把「全异步 RL」讲明白

从一张四宫格图出发,讲清楚 on-policy、流式 off-policy、异步陈旧样本、Partial Rollout 这四种训练模式到底差在哪;再顺着一条样本的旅程,走完 Rollouter → Harbor 沙箱 → MessageQueue → Trainer → 权重同步的完整链路。

素材:四种RL模式对比.pngVerl+harbor全异步RL分享.md
读者:了解 PPO/GRPO 基本概念,但没读过 VeRL 全异步源码的人

01先说结论:这套东西在解决什么

在 SWE-bench 这类代码 Agent 场景里做 RL,一条样本要在沙箱里跑几分钟到几十分钟,而且每条样本的耗时天差地别。传统训推一体(Colocate)架构下,GPU 会长时间空转等最慢的那条样本。全异步(Fully Async)把生成和训练拆到两组 GPU 上、用一个消息队列以「单条样本」为粒度流式对接,端到端吞吐提升 2.35×–2.67×

名词解释 · 先扫清这几个词
SWE-bench
用真实 GitHub 仓库的 issue 出题、让模型改代码、跑测试判对错的基准。一条题目就是一次「修 bug 任务」。
代码 Agent
不是一问一答,而是模型在环境里反复「思考 → 调工具(bash、编辑器)→ 看结果」的多轮循环,直到把任务做完。
沙箱(Sandbox)
隔离的运行环境(这里是 K8s Pod 或 Docker 容器)。Agent 在里面随便执行命令、改文件,都不影响外面的训练系统。
Rollout / 轨迹
让当前模型实际跑一遍任务、产出一条完整交互记录(多轮对话 + 工具调用 + 最终结果),这条记录就是一条训练样本。
Colocate(训推一体)
生成样本(推理)和更新参数(训练)共用同一组 GPU,只能轮流干活。
吞吐(Throughput)
单位时间能完成多少样本/训练步,衡量整套系统「快不快」。

而那张四宫格图,画的是同一套架构下的四档「异步程度」:从最保守的完全同步(a),一路放宽到最激进的「同步权重时把生成中的样本切断、同步完再接着写」(d)。越往后吞吐越高,代价是训练用的样本越「不新鲜」,算法上要付出额外的修正手段。

02痛点:GPU 一半时间在发呆

Colocate 模式下,Rollout(生成轨迹)和 Training(更新参数)共用同一组卡,串行执行:

Colocate:
  |------ rollout 1(含长尾等待)------|-- train --|------ rollout 2 ------|-- train --|
                    ↑
        绝大多数样本早就跑完了,
        整批在等最慢的那 1~2 条
打个比方

一个厨师既炒菜又洗碗,而且规定「一桌 64 道菜全部出锅才能开始洗碗」。有 62 道菜 3 分钟就好了,剩下 2 道要炖 40 分钟——于是灶台空着 37 分钟,锅碗也堆着不洗。这就是长尾样本 + 训推串行的双重浪费。

全异步的做法是:请两个厨师,各占一个灶台——一个只管炒菜(Rollouter),一个只管洗碗(Trainer),中间放一个传菜口(MessageQueue)。菜一道一道地传,不必等整桌齐:

Fully Async:
  Rollouter: |--- 持续生成 ---|--- 持续生成 ---|--- 持续生成 ---|
  Trainer:            |---- train 1 ----|---- train 2 ----|---- train 3 ----|
             ══════════════> 两条时间线重叠,端到端时间大幅压缩

03四大组件:把厨房拆成两半

┌──────────────────────────────────────────────────────────────┐
│                     Fully Async Policy                        │
│                                                               │
│   ┌───────────┐   一条一条传    ┌──────────────┐              │
│   │ Rollouter │ ─────────────> │ MessageQueue │              │
│   │  (8 GPU)  │                │  (Ray Actor) │              │
│   └───────────┘                └──────┬───────┘              │
│        ▲                              │ 一条一条取            │
│        │ NCCL 广播权重                 ▼                      │
│   ┌────┴──────────────────┐    ┌──────────────┐              │
│   │ ParameterSynchronizer │ <──│   Trainer    │              │
│   │  (checkpoint-engine)  │    │   (8 GPU)    │              │
│   └───────────────────────┘    └──────────────┘              │
└──────────────────────────────────────────────────────────────┘
组件身份做什么
Rollouter生产者独占一组 GPU(如 8 卡),不停地跑 Agent 轨迹,每完成一条就立刻塞进队列。生成速度由 staleness 参数节流。
MessageQueue缓冲区一个 Ray Actor,内部就是 deque + asyncio.Condition 的生产者-消费者队列。
Trainer消费者独占另一组 GPU,从队列里逐条取样本,攒够 require_batches × ppo_mini_batch_size 条就训一步。
ParameterSynchronizer搬运工基于 NCCL(参考 checkpoint-engine),把 Trainer 的新权重高效广播给 Rollouter。
名词解释 · 这张组件图里的基础设施
Ray / Ray Actor
Ray 是 Python 的分布式计算框架。Actor 是一个常驻的远程服务对象:起在集群某个进程里,别的进程可以远程调用它的方法。MessageQueue、LoadBalancer 都是 Ray Actor——所以全集群共享同一份状态。
deque
Python 内置的双端队列,这里当先进先出的缓冲区用。
asyncio / 协程
Python 的单线程异步并发:成百上千个任务在等 IO(网络、推理返回)时互相让出 CPU,谁的数据到了谁继续跑。asyncio.Condition 就是配套的「队列空了先睡、来货了叫醒」的通知机制。
生产者-消费者
经典并发模式:一方只管往队列塞(Rollouter),一方只管从队列取(Trainer),两边速度不匹配由队列缓冲吸收。
NCCL
NVIDIA 的 GPU 集合通信库,让多卡/多机的显存之间直接高速传数据(不绕 CPU)。分布式训练里同步权重、梯度的标配,比走普通网络快一个量级。
关键设计:传输粒度是「一条样本」

不是一个 batch,是单条轨迹。这一点是整个流式架构的地基——只有粒度足够细,Rollouter 才不用等齐一整批,Trainer 也才能边攒边算。

04读懂那张四宫格图

四种 RL 模式对比全图
原图全貌。四个子图共用同一套坐标:横轴是时间,纵轴从上到下是 Rollouter / MessageQueue / Trainer 三条泳道。

先花 30 秒把「怎么看」搞清楚,后面四张图就都能秒懂了。

图例

三条泳道

  • Rollouter(蓝条):每一根横条 = 一条正在生成的轨迹。条子长度不一,画的就是长尾——有的样本几分钟,有的几十分钟。
  • MessageQueue(竖格):队列里堆积的样本,一格一条。
  • Trainer(绿框):每个 mini-batch 方块 = 一次本地参数更新。

颜色的含义

  • 黄色竖线 = 参数同步。这是全图的节拍器:黄线左右两侧,Rollouter 用的是不同版本的模型权重。
  • 灰格 = 新鲜样本:用当前权重版本生成的。
  • 粉格 = 陈旧样本(staleness):用上一个权重版本生成、却拿到这一轮来训练的样本。
  • 半灰半粉 = partial 样本:一条轨迹横跨了权重同步点——前半段是旧权重生成的,后半段是新权重接着生成的。
  • 子图 (d) 里 Rollouter 泳道上那些橙色虚线条,就是被中断、然后在下一段接着写完的轨迹。

红框注解 = 在等谁

wait first batch rollout finishwait last batch train finishwait active task finish——这些红框标的都是气泡(bubble),也就是这个模式无法消除的空转。看图时最该盯的就是红框:红框越少,流水线越满。

05四种模式逐格拆解

四种模式由三个开关组合而成,先把开关摆出来:

模式trigger_parameter_sync_stepstaleness_thresholdpartial_rollout一句话
(a) On Policy= 10训一步同步一次,最保守
(b) Stream Off Policy> 10训多步才同步一次
(c) Async (stale)≥ 1> 0False允许旧样本,但同步时要等在跑的样本跑完
(d) Async (partial)≥ 1> 0True允许旧样本,同步时直接把在跑的样本切断、之后续写
(a)

On Policy Pipeline —— 最干净,也最慢

模式 a:on policy pipeline
图上特征
黄线密集;每两条黄线之间,Rollouter 只有寥寥几条蓝条,Trainer 只有 1 个 mini-batch;队列格子全是灰的。
节奏
生成一批 → 训一步 → 立刻同步权重 → 再生成一批。
代价
黄线两侧都是空白:同步期间 Rollouter 停着,训练期间队列是空的。

这是严格 on-policy:训练用的每一条样本,都是当前这版权重刚生成出来的。算法最干净,PPO 的重要性采样比值天然接近 1,不需要任何修正。但代价是流水线几乎没有重叠——图上黄线之间的空白就是纯浪费。

名词解释 · on-policy / off-policy / 重要性采样
on-policy
训练用的样本必须是「当前这版模型」自己刚生成的——采样的策略和被优化的策略是同一个。
off-policy
样本来自旧版本模型(或别的策略)。数据分布和当前策略对不上,直接拿来训会有偏差,需要修正。
重要性采样(IS)
修正的手段:给每个 token 乘一个比值 π_new / π_old(新旧策略给这个 token 的概率之比),把「旧分布下采到的样本」换算到新分布。比值越接近 1,说明新旧策略差别越小、修正越轻——这就是正文说 (a) 「比值天然接近 1」的意思。
mini-batch
一次参数更新实际用到的那撮样本。Trainer 每攒够 require_batches × ppo_mini_batch_size 条就做一次本地更新——图上 Trainer 泳道的一个绿方块就是一个 mini-batch。
(b)

Stream Off Policy —— 少同步几次,但两头还是要等

模式 b:stream off policy pipeline
图上特征
黄线稀疏了;每个窗口里 Trainer 连做 3 个 mini-batch;队列格子仍然全是灰的;但出现了两个红框。
红框 1
wait first batch rollout finish —— 窗口开头,Trainer 得干等第一批样本攒够。
红框 2
wait last batch train finish —— 窗口结尾,Rollouter 已经把额度用完停下了,得等 Trainer 训完最后一步才能同步。

trigger_parameter_sync_step > 1 的意思是:Trainer 本地更新 N 次,才和 Rollouter 同步一次权重。同步是有固定开销的(abort、释放 KV cache、建 NCCL 通信组、广播、唤醒),少同步几次自然更省。

但注意:staleness 仍然是 0。也就是说 Rollouter 在一个窗口内只被允许生成固定数量的样本,生成完就得停下等同步。所以图上窗口的头和尾各有一个气泡——这正是 (c) 要解决的问题。

为什么 (b) 已经算 off-policy 了

因为一个窗口里训了 3 步。第 2、3 步用的样本,是窗口开始时那版权重生成的,而此时 Trainer 的参数已经更新过 1~2 次了。样本和参数已经对不齐——这就是 off-policy 的来源,只不过偏差还比较可控。

(c)

Async Stream + 陈旧样本 —— 让生成永不停机

模式 c:async stream pipeline with staleness samples
图上特征
Rollouter 泳道密密麻麻铺满,几乎没有空白;队列里出现了粉色格子;红框只剩一个:wait active task finish
核心变化
staleness_threshold > 0:允许 Rollouter「超额生产」,多出来的样本留给下一个窗口用。
剩余气泡
要同步权重时,必须等所有正在跑的样本跑完才能动权重——长尾又回来了,只不过只在同步点上出现一次。

关键洞察是:Rollouter 停机的根本原因,是「本窗口配额用完了」。那就把配额放宽——允许它提前生成下一窗口要用的样本。这些样本进入队列时是新鲜的(灰色),但被 Trainer 取走时权重已经换了一版,于是变成陈旧样本(粉色)。

打个比方

厨师照着上一版菜谱已经把菜炒好了,端上桌时菜谱刚好更新了。菜还是能吃的,只是不完全符合新标准。staleness_threshold 就是在规定「一桌菜里,最多允许多大比例是照旧菜谱做的」。

这就是全异步真正开始赚钱的地方:Rollouter 的空转基本被填平了。代价是训练数据里混入了 off-policy 成分,需要在算法侧留心(见 §06 的 bypass_mode 与 Rollout IS)。

(d)

Async Stream + Partial Rollout —— 把最后一个气泡也挤掉

模式 d:async stream pipeline with partial rollout
图上特征
一个红框都没有了。Rollouter 泳道上出现橙色虚线条——被切断又续上的轨迹;队列里出现半灰半粉的格子。
核心变化
partial_rollout=True:同步权重时不再等待,直接 abort 所有在途请求;同步完成后 resume,从中断点接着生成。

(c) 里最后那个气泡是「等长尾样本收尾」。而一条 SWE 轨迹可能要跑几十分钟,为它等待代价太大。Partial Rollout 的答案很直接:不等了,切断它

打个比方

你在写一篇长文,写到一半编辑部换了写作规范。你不用把整篇按旧规范写完再重写,而是停笔、换上新规范、从当前这句接着往下写。文章前半段是旧规范、后半段是新规范——这就是那个「半灰半粉」的格子。

实现上有两个漂亮的细节:

名词解释 · KV Cache 一家子
KV Cache
Transformer 生成时缓存的中间结果:每个已处理 token 的 Key/Value 向量。有了它,生成下一个 token 时不用把前文从头重算一遍。它存在 GPU 显存里、跟着具体某台推理实例走。
prefill
正式生成前,把整段 prompt 过一遍模型、填好 KV Cache 的阶段。prompt 越长 prefill 越贵——多轮 Agent 对话动辄几万 token,prefill 是大头。
APC
vLLM 的 Automatic Prefix Caching(自动前缀缓存):两次请求只要有相同的前缀,前缀部分的 KV Cache 直接复用、不再 prefill。多轮对话里「前几轮的全部内容」正好就是公共前缀,所以收益巨大——前提是请求得落在同一台实例上。
request_id
VeRL 里一条轨迹(整段多轮对话)的身份证号。续写认它、负载均衡的粘性路由也认它——所以中断重试能回到原来那台实例、接上原来的缓存。

轨迹会记录 min_global_steps / max_global_steps,标明它横跨了哪几个权重版本,供算法侧做修正。

四格横向对比

Rollouter 空转Trainer 空转样本新鲜度适合什么场景
(a)100% 新鲜算法验证、小规模、对 off-policy 极敏感的任务
(b)中(首尾各一个气泡)100% 新鲜同步开销显著、rollout 耗时较均匀
(c)小(只剩同步点等长尾)混入陈旧样本rollout 耗时中等、长尾不算极端
(d)≈ 0≈ 0陈旧 + 跨版本轨迹代码 Agent 这类超长尾场景
一句话记住这张图

从 (a) 到 (d) 是逐个消灭气泡的过程:(b) 消灭「频繁同步的开销」,(c) 消灭「Rollouter 等配额」,(d) 消灭「同步时等长尾」。而每消灭一个气泡,就多欠算法一点 off-policy 的债——所以 staleness_threshold 建议设成小于 1 的值,别把债欠太狠。

06三个旋钮:参数怎么调

trigger_parameter_sync_step:多久同步一次权重

Trainer 做多少次本地更新后,才和 Rollouter 同步一次。两次同步之间,Trainer 一共消费:

trigger_parameter_sync_step × require_batches × ppo_mini_batch_size  条样本

想和 colocate 模式做公平的速度对比,就把它设成:

trigger_parameter_sync_step = data.train_batch_size / (require_batches × ppo_mini_batch_size)

staleness_threshold:允许多少「不新鲜」

它是允许使用的过期样本最大比例。两次权重同步之间,Rollouter 最多生成:

staleness_threshold = 0(同步训练):
    rollout_num = trigger_parameter_sync_step × require_batches × ppo_mini_batch_size

staleness_threshold > 0(异步训练):
    rollout_num = (1 + staleness_threshold)
                × (trigger_parameter_sync_step × require_batches × ppo_mini_batch_size)
                − num_staleness_sample

其中 num_staleness_sample 是上一轮多生产、结转过来的样本数——所以这是个自动结算的额度机制:上轮超产多少,这轮就少产多少,总量守恒。

怎么选值

partial_rollout:只在 staleness > 0 时才真正生效

这一点容易踩坑:staleness_threshold = 0 时打开 partial_rollout 是没有意义的——因为窗口内本来就不允许有跨版本样本存在。

require_batches:一次训练喂多少

理论上流式训练该设成 1(攒够一个 ppo_mini_batch_size 就训)。但实测发现:一次分发的样本太少,会因为数据分发的顺序问题导致训练不稳定、response 长度变长require_batches 就是为此提供的一个调节量,用来控制每次真正参与训练的样本数。

use_rollout_log_probsbypass_mode:一个容易被忽略的正确性问题

PPO/GRPO/DAPO 计算重要性采样时,old_log_prob 必须和「生成这些 token 时所用的那一版参数」对应。在全异步模式下,样本可能是好几个版本以前生成的——如果让 Trainer 用当前权重去重算 log_prob,比值就错了。

bypass_mode = True(默认)
    old_log_probs = rollout 侧返回的 log_probs   ← 零额外计算,天然版本对齐

bypass_mode = False
    old_log_probs = 用训练引擎(FSDP/Megatron)重新计算
    同时启用 Rollout Importance Sampling 修正
    在模式 (d) 下,这近似于 AReaL 的 Decoupled PPO

什么时候关掉 bypass?文档给的经验是:训练后期指标和 response 长度开始不稳定时,改用训练引擎算 log_prob 并叠加 Rollout IS 修正,能缓解这个问题。

名词解释 · log_prob 与几个算法名词
log_prob
模型给「自己生成的每个 token」打的对数概率。PPO 用它来算重要性采样比值(见 §05)。
old_log_prob
比值的分母,代表「生成这条样本的那个策略」。关键约束:它必须来自生成时那一版参数——用错版本,整个比值就失去意义。
one-step-off
一种常见的轻度异步方案:训练用的样本恰好、且始终落后一个参数版本。staleness_threshold = 1 且 rollout 足够快时,系统会自然收敛到这个稳态。
Rollout IS / Decoupled PPO
把「实际生成样本的策略」和「PPO 比值里的行为策略」解耦,再用重要性采样把差距补回来的一类方法;AReaL 论文把它叫 Decoupled PPO。bypass_mode=False 走的就是这条路。
FSDP / Megatron
两种主流的分布式训练引擎(模型切分方式不同)。这里只需知道:它们是「训练侧」的模型,与「推理侧」的 vLLM 相对。

07一条样本的旅程:六层并发漏斗

从 Rollouter 拿到一条任务,到它在 Harbor 沙箱里跑完,中间要穿过六层并发限制。以 fully_async_2nodes.sh 的配置为例(rollout 8 GPU、gen_tp=4num_workers=128ppo_mini_batch_size=64staleness_threshold=1.0):

max_required_samples   = 64 × (1.0 + 1) × 1        = 128   # 一个同步窗口内最多生成
num_replicas           = 8 / 4                      = 2     # vLLM server 实例数
max_concurrent_samples = min(2 × 32, 128)           = 64    # 同时在途的样本数上限
max_queue_size         = max_required_samples       = 128
Dataloader — 逐条读取任务
① pending_queue · 上限 128
有界 asyncio.Queue,满了就阻塞 feed 协程,防止把数据全读进内存
② active_tasks · 上限 64
Rollouter 侧唯一的硬限制。满了就 asyncio.wait(FIRST_COMPLETED) 等一个完成再提交新的
③ AgentLoopWorker 池 · 128 个 Ray Actor
上游只有 64 在途,所以这层当前配置下不是瓶颈
④ 单 Worker 内 asyncio.gather
一个 Worker 收到 batch 后把每条样本各起一个 Task 并行跑
⑤ Harbor 沙箱 · 无软件限制
K8s Pod 数量只受集群 CPU/内存配额限制;资源不够就排队等调度
⑥ vLLM Server · 2 个 replica(dp2 tp4)
所有 LLM 推理的物理上限。好在样本是多轮交互、节奏交错,并非同时压上来
瓶颈排序(从紧到松)

vLLM replica(硬件)→ active_tasks = 64(软件硬限)→ pending_queue = 128(缓冲)→ Worker 池 128(容量)→ K8s 集群资源(基础设施)。
换句话说:想提高并发,先看 GPU 和 K8s 有没有余量,再动 max_concurrent_samples

名词解释 · 漏斗里出现的词
有界队列
设了容量上限的队列,满了就塞不进去——上游被迫等待。这是最简单可靠的限流手段,防止把整个数据集一口气读进内存。
replica(副本)
一个完整、可独立提供服务的 vLLM 实例。这里 8 张卡切成 2 个 replica,每个各占 4 张卡。
dp / tp
两种并行方式:tp(张量并行)把一份模型切到几张卡上合力算一个请求;dp(数据并行)把模型复制几份、各自接不同请求。dp2 tp4 = 2 个副本 × 每副本 4 卡张量并行。
asyncio.gather
把一批协程任务同时挂起来并行跑、等它们全部完成。一个 Worker 用它同时驱动手头的多条轨迹。

GlobalRequestLoadBalancer:粘性会话 + 最少在途

所有 AgentLoopWorker 共享一个全局 Ray Actor 做路由,规则只有两条:

  1. 粘性会话:同一个 request_id(= 同一条轨迹)的所有轮次,通过 LRU Cache 记录,始终打到同一个 replica。目的是让 vLLM 的 APC 复用前几轮的 KV Cache,省掉大量 prefill。
  2. 最少在途请求:新对话分配给当前 inflight 计数最小的 replica,避免一个忙死一个闲着。

请求结束时在 finally 里 fire-and-forget 地 release_server,把计数减回去。Partial Rollout 的续写之所以能复用缓存,正是因为它沿用了同一个 request_id,被粘性会话送回了原来那台 replica。

名词解释 · 负载均衡这几个词
LRU
Least Recently Used,「最近最少使用」淘汰策略:缓存放满时,先扔掉最久没被访问的那条。活跃对话每来一次请求都会刷新自己的位置,不会被挤掉;被挤掉的多半是早已结束的对话,丢了也无所谓。
LRU Cache
按 LRU 策略封顶的字典。这里存的是路由表 request_id → replica(上限 10000 条)——所以「命中 LRU」= 这条对话在路由表里还有记录,走粘性路径;「未命中」= 新对话或记录已被淘汰,重新挑一台。
粘性会话(Sticky Session)
Web 后端的经典概念:同一个会话的所有请求,固定发给同一台服务器。这里「粘」的目的只有一个——让 APC 能复用前几轮的 KV Cache。
在途请求(inflight)
已经发出、还没返回的请求数。它是衡量一台 replica 忙闲的实时指标,「最少在途」就是把新对话派给最闲的那台。
fire-and-forget
调用发出去就不等结果返回。release 只是把计数减一,这种小事不值得让 Worker 阻塞着等一次远程调用。

08Harbor 适配与那个「假装是 vLLM」的 Proxy

BuiltinSWEAgentLoop:把整个 Harbor Trial 塞进一次 run()

VeRL 原本的 ToolAgentLoop 是自己驱动状态机、自己在进程内执行工具的。但 SWE 场景下,Agent(如 OpenHands)跑在 K8s Pod 里,循环是 Harbor 控制的。所以适配的思路是:把一整次 Harbor Trial.run() 包装成 VeRL 的一次 run() 调用

ToolAgentLoopBuiltinSWEAgentLoop
谁控制 Agent 循环VeRL(状态机驱动)Harbor(Trial.run 内部)
工具在哪执行VeRL 进程内Harbor 沙箱(K8s Pod / Docker)
怎么调模型直接 server_manager.generate()经 HTTP Proxy 间接触发
轨迹 token 哪来状态机自己累积Proxy 会话捕获
怎么算 rewardVeRL RewardModelHarbor Verifier(真跑测试套件,0/1)

一次 Trial 内部分四个阶段,每个阶段都有独立计时并上报到 VeRL metrics:env_setup(起 Pod、跑 Dockerfile)→ agent_setup(装 Agent、把 LLM_BASE_URL 指向 Proxy)→ agent_execute(多轮「推理 ↔ 工具」循环)→ verify(跑测试算 reward)。

Proxy:一个进程内的假 vLLM 端点

沙箱里的 Agent 用的是标准 LiteLLM/OpenAI 客户端,它不知道 VeRL 的存在。于是 VeRL 在自己进程里起了一个 aiohttp server,伪装成 OpenAI 兼容的 vLLM 端点,把 LLM_BASE_URL 塞给 Agent。Agent 一行代码都不用改。

Harbor Agent (K8s Pod)          Proxy (进程内 aiohttp)          VeRL
      │  POST /sess/{id}/v1/chat/completions │                    │
      │ ───────────────────────────────────> │                    │
      │                    ① 归一化 messages  │                    │
      │                    ② apply_chat_template                  │
      │                    ③ 算出 prompt_ids  │                    │
      │                                      │ server_manager     │
      │                                      │  .generate()  ───> │
      │                                      │ <─── TokenOutput   │
      │                    ④ 解析 tool_call(hermes / qwen3_coder)│
      │                    ⑤ 更新 session 轨迹状态                 │
      │ <──── OpenAI 格式 JSON ───────────── │                    │

它做的事远不止转发——它是轨迹 token 的唯一采集点。因为 Agent 在沙箱里,VeRL 看不到对话历史,只能在 Proxy 这一层把每一轮的 token、mask、logprob 攒起来。

名词解释 · Proxy 相关
OpenAI 兼容端点
实现了 OpenAI API 格式(POST /v1/chat/completions,收 messages、返回 choices)的 HTTP 服务。它是 LLM 服务的事实标准——任何 OpenAI SDK / LiteLLM 客户端都能直连。vLLM 本身就长这样,所以 Proxy「装成 vLLM」毫无违和感。
LiteLLM
把各家 LLM API 统一成一套接口的 Python 客户端库,Agent 框架常用它对接任意模型服务。沙箱里的 Agent 用的就是它。
aiohttp
Python 的异步 HTTP 库。这里用它在 VeRL 进程内真正监听一个端口、当 HTTP 服务器用——因为它和 VeRL 共用同一条 asyncio 事件循环,几十条并发轨迹的请求可以同时挂起等推理,互不堵塞。

Token-in-Token-out:整个适配层最硬的那块骨头

Trainer 需要精确知道每个 token 的身份(哪些是模型生成的、哪些是工具返回的)。Proxy 是增量累积的——每轮只 tokenize 新增的那几条消息。但这个增量结果,必须和「把完整对话一次性 apply_chat_template」的结果逐 token 完全一致,否则训练数据就是错位的。

打个比方

把 10 样东西一件一件称重再加起来,必须和一次性称这 10 样东西的读数一模一样。听起来是废话,但 chat template 里有各种分隔符、系统提示、生成前缀,稍不留神就多一个或少一个 token。

Proxy 用三招守住这个不变式:

此外 Proxy 还负责把 log_probs(作为 rollout_log_probs)、response_maskrouted_experts(MoE 的 Router Replay 用)一路传回 Trainer;并支持 max_consecutive_no_tool 早停,防止 Agent 陷入「只聊天不干活」的死循环。

名词解释 · token 加工链上的词
tokenize
把文本切成 token 并映射成整数 id。模型和 Trainer 只认 token id,不认字符串。
chat template
把 messages 列表(system/user/assistant/tool)拼装成模型真正吃的那串文本的模板:加 <|im_start|> 之类的角色分隔符、工具描述、生成前缀。每家模型的模板都不一样,token 错位的坑大多埋在这里。
EOS
end-of-sequence,模型表示「这轮我说完了」的特殊 token。vLLM 生成到 EOS 就停,但模板里 EOS 后面可能还有字符(比如一个换行)——这就是「EOS 尾部回填」要补的东西。
response_mask
和轨迹等长的 0/1 序列:1 = 模型自己生成的 token(参与算 loss),0 = 工具输出、observation 等外部塞进来的 token(不算 loss)。RL 只应该奖惩模型自己说过的话。
hermes / qwen3_coder
两种 tool-call 的文本格式约定(及其解析器)。模型把「我要调用哪个工具、参数是什么」写在输出文本里,解析器把它还原成结构化的 tool_calls 字段。
MoE / routed_experts
MoE(混合专家)模型里,每个 token 只会被路由到少数几个「专家」子网络。routed_experts 记录了生成时每个 token 实际走了哪些专家,训练时按记录回放(Router Replay),保证训推一致。

09训练侧:Trainer 每一步干什么

FullyAsyncTrainer.fit() 是个无限循环,每轮 fit_step() 十个动作:

 1. _fit_generate           从 MessageQueue 逐条取样本,攒够 required_samples 组装成 batch
 2. _fit_compute_reward     reward 评分(Harbor 侧其实已经算好了)
 3. _fit_compute_log_prob   old_log_prob(bypass_mode 决定是直接用还是重算)
 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       PPO / GRPO / DAPO loss 更新 Actor
 9. _fit_update_local_step  local_trigger_step += 1,到达阈值则复位
10. _fit_update_weights     若 local_trigger_step == 1:
                              → checkpoint_manager.update_weights()   同步权重
                              → rollouter.reset_staleness()           重置陈旧度基准

第 9、10 步就是「多步本地训练 + 一次参数同步」的实现:计数器从 1 数到 trigger_parameter_sync_step,然后复位并把 current_param_version 加一。

名词解释 · 训练步里的角色
Actor
被训练的策略模型本体——就是那个生成代码、调工具的 LLM。PPO 里「更新 Actor」就是更新它的参数。
Critic
价值网络:估计「从当前状态出发预期能拿多少 reward」,帮 Actor 判断某一步比平均好多少。GRPO 不用 Critic,PPO 用。
Reference Policy
一个冻结不动的参考模型(通常是训练起点的 SFT 模型)。用 KL 散度惩罚拽住 Actor,别为了刷 reward 跑到胡言乱语的分布上去。
GAE
Generalized Advantage Estimation:把整条轨迹的 reward 折算成「每一步的优势值」(这步比预期好多少)的标准算法,比直接用总回报方差更低、训练更稳。

10权重同步与断点续传

CheckpointEngineManager.update_weights() 的八个步骤:

#动作为什么
1abort_all_requests()切断所有在途推理,返回 stop_reason="aborted"这是 partial rollout 的起点
2构建临时 RayWorkerGroup把所有 rollout replica 的 worker 收拢起来
3sleep replicas(可选)释放 vLLM 的 KV-cache 显存,给 NCCL 传输腾地方
4建 NCCL 通信组Trainer rank 0 + 全部 rollout worker
5权重传输Trainer 侧 pack 进 CuPy 双缓冲后 NCCL Broadcast;元数据走 ZeroMQ PUB/SUB 并行发送,与 NCCL 传输重叠
6Finalize释放 NCCL bucket、销毁通信组
7wake up replicas恢复 KV-cache,新权重上 GPU
8resume_generation()被中断的样本从断点继续生成——这就是 (d) 里那些橙色虚线条
名词解释 · 权重搬运用到的工具
Broadcast
集合通信的基本操作之一:一个节点(Trainer rank 0)把同一份数据发给通信组里所有其他节点(全部 rollout worker)。
CuPy 双缓冲
CuPy 是「GPU 版 numpy」。双缓冲 = 准备两块显存轮流用:一块正在 NCCL 发送时,另一块同时打包下一批参数——传输和打包重叠,不互相等。
ZeroMQ PUB/SUB
轻量消息库的「发布/订阅」模式。大块权重走 NCCL,而参数的元数据(名字、形状、走到哪了)走 ZeroMQ 并行发送,两条通道互不干扰。
sleep / wake up
vLLM 提供的接口:sleep 释放 KV Cache 占用的显存(给 NCCL 传输腾地方),wake up 再重建。这就是表中第 3、7 步做的事。

Partial Rollout 的实现只有一个 while 循环

class FullyAsyncLLMServerManager:
    async def generate(self, ...):
        final_output = TokenOutput(token_ids=[], ...)
        while True:
            # 每次都把 prompt + 已生成的 token 一起传进去
            output = await super().generate(
                prompt_ids=prompt_ids + final_output.token_ids, ...
            )
            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:
                continue      # 被权重同步打断了 → 接着写
            else:
                break         # 正常结束
        return final_output

上层的 Agent Loop 完全感知不到这中间发生了权重切换——它只看到一次普通的 generate()

reset_staleness:同步之后怎么重新记账

async def reset_staleness(self):
    async with self.lock:
        # 在途任务 + 队列里等待的样本,都算作新版本下的样本
        self.staleness_samples = len(self.active_tasks) + self.pending_queue.qsize()
        self.paused = False
        self._resume_event.set()

逻辑是:staleness_samples 要反映「当前参数版本之后已经生成了多少」。同步完成的那一刻,手上那些还没交付的样本,从记账角度就归到新版本名下——下一个窗口的生成配额会自动扣掉它们。这就是 §06 里那个「总量守恒」的实现。

11术语表

架构与模式

术语解释
Colocate训推一体:Rollout 和 Training 共享同一组 GPU,只能轮流执行
Fully Async全异步训推分离:Rollouter 和 Trainer 各占一组 GPU、完全解耦,靠队列流式对接
Rollouter样本生成器:持续异步产出训练样本,每完成一条立刻入队
MessageQueue基于 Ray Actor 的异步队列,以「单条样本」为粒度连接 Rollouter 和 Trainer
Rollout / 轨迹让当前模型实际跑一遍任务产出的完整交互记录(多轮对话 + 工具调用),即一条训练样本
Staleness(陈旧度)样本生成时的参数版本与当前训练参数版本的差距;staleness_threshold 限制陈旧样本的最大比例
Partial Rollout中断恢复:同步权重时打断进行中的生成,同步后带着已生成的 token 从断点续写
on-policy / off-policy训练样本是否由当前这版策略自己生成;不是,就是 off-policy,需要重要性采样等修正
重要性采样(IS)用概率比值 π_new/π_old 把旧策略采到的样本换算到新策略分布下的修正方法
MISMultiple Importance Sampling,多步重要性采样:样本横跨多个参数版本时保证 π_old 一致性
one-step-off训练样本恰好、且始终落后一个参数版本的稳态异步方案
bypass_mode直接用 rollout 侧返回的 log_probs 当 old_log_probs——零额外计算且天然版本对齐
气泡(bubble)流水线里无法避免的 GPU 空转时段,即四宫格图里的红框

推理与缓存

术语解释
vLLM主流的高吞吐 LLM 推理引擎,对外提供 OpenAI 兼容的 HTTP 端点
replica(副本)一个完整、可独立服务的 vLLM 实例;dp2 tp4 = 2 个副本 × 每副本 4 卡张量并行
KV Cache推理时缓存的每个已处理 token 的 Key/Value 向量,生成新 token 时免去重算前文
prefill生成前把整段 prompt 过一遍模型、填好 KV Cache 的阶段,prompt 越长越贵
APCvLLM 的 Automatic Prefix Caching:相同前缀的请求直接复用已有 KV Cache,多轮对话复用的关键
LRULeast Recently Used(最近最少使用):缓存满时先淘汰最久没被访问的条目
LRU Cache按 LRU 策略封顶的字典;LoadBalancer 用它存 request_id → replica 路由表(上限 10000)
粘性会话Sticky Session:同一会话的所有请求固定路由到同一台服务器,为的是复用 KV Cache
在途请求(inflight)已发出、尚未返回的请求数,衡量一台 replica 忙闲的实时指标
request_id一条轨迹(整段多轮对话)的唯一标识,粘性路由与断点续写都以它为准

Harbor 适配与 Proxy

术语解释
Proxy进程内 aiohttp HTTP Server,伪装成 vLLM 端点,把 Harbor 的 OpenAI 请求转成 VeRL 的 server_manager 调用,并顺路采集轨迹 token
AgentLoopVeRL 中封装「生成单条轨迹」逻辑的抽象单元
OpenAI 兼容端点实现 OpenAI API 格式(/v1/chat/completions)的 HTTP 服务,LLM 服务的事实标准接口
LiteLLM统一各家 LLM API 的 Python 客户端库,Agent 框架对接模型服务的常用选择
aiohttpPython 异步 HTTP 库,与 VeRL 共用一条 asyncio 事件循环,高并发请求互不堵塞
chat template把 messages 列表拼成模型实际输入文本的模板(角色分隔符、工具描述、生成前缀)
EOSend-of-sequence,模型表示「说完了」的特殊 token;EOS 之后的模板尾巴需要 Proxy 手动回填
response_mask0/1 序列:1 = 模型生成的 token(算 loss),0 = 工具返回等外部 token(不算)
Token-in-Token-out增量 tokenization 与一次性 tokenization 的结果必须逐 token 一致的不变式
Router ReplayMoE 路由决策的复现:记录生成时每个 token 走过的专家(routed_experts),训练时按记录回放

基础设施

术语解释
Ray / Ray ActorRay 是分布式计算框架;Actor 是常驻的远程服务对象,全集群可远程调用、共享其状态
asyncio / 协程Python 单线程异步并发:大量任务在等 IO 时互相让出 CPU
NCCLNVIDIA 的 GPU 集合通信库,多卡/多机显存间直接高速传输,权重同步的标配
Broadcast集合通信操作:一个节点把同一份数据发给通信组内所有其他节点
checkpoint-engine基于 NCCL 的高效分布式参数同步引擎(本架构权重同步的参考实现)
CuPyGPU 版 numpy;权重同步时用它做双缓冲打包,让打包与传输重叠
ZeroMQ轻量消息库;PUB/SUB 模式用于并行发送权重元数据,与 NCCL 大块传输互不干扰
fire-and-forget调用发出后不等待结果返回,适合无关紧要的收尾操作
FSDP / Megatron两种主流分布式训练引擎,即「训练侧」的模型载体,与推理侧的 vLLM 相对

回到那张图

如果只带走一句话:四宫格画的是同一条流水线在四种「松紧度」下的样子,红框标的就是浪费掉的 GPU 时间,而 (d) 是唯一一张没有红框的图。代价写在队列那一行——越往后,粉色和半粉的格子越多,off-policy 的债就越重。staleness_threshold 是你在吞吐和精度之间的那个刻度盘。

← 全部解读