首页
学习
活动
专区
圈层
工具
发布
社区首页 >专栏 >29-FullyAsyncPolicy速度新鲜度和off-policy的三角关系

29-FullyAsyncPolicy速度新鲜度和off-policy的三角关系

作者头像
anzhsoft
发布2026-07-27 18:48:24
发布2026-07-27 18:48:24
1330
举报

Fully Async Policy 的核心不是“rollout 更快”,而是把 trainer 等 rollout 的同步边界改成 producer/consumer 系统,然后把样本新鲜度变成必须度量和约束的状态。 当 trainer 不再等待当前 step 的 rollout,而是从持续生产的队列中凑 batch,吞吐会提升,样本新鲜度也变成 系统风险。本文从 MessageQueue、staleness、partial rollout 到 rollout correction,讲清异步训练的三角账

第 28 篇讲 TransferQueue:它把 rollout 产物从 controller 返回值改成外部队列和 metadata,但 trainer 仍然围绕当前 step 等样本。第 29 篇继续往前走:如果 rollouter 可以持续生产,trainer 可以从队列里凑 batch,系统吞吐会提升,但训练语义也会改变。

本文的核心判断是:Fully Async Policy 用 MessageQueue把 trainer 和 rollouter 拆成两个长期运行的 actor;trainer 按 required_samples从队列取样本,rollouter 按并发数、队列大小和 staleness 上限持续生成。它买到的是 rollout 与训练 update 的重叠,付出的账是样本可能来自旧参数版本,必须通过 staleness、partial rollout、parameter sync 和 rollout correction 控制 off-policy 风险。

先看全局图。读图时注意三条线:样本线从 rollouter 进 queue 再到 trainer;参数线从 trainer 同步回 rollout server;诊断线围绕 queue、staleness 和 correction 指标闭环。

Fully Async Policy 的三角关系

这张图补上第 28 篇没有跨过去的边界:TransferQueue 解决的是数据怎么存、怎么按字段读写;Fully Async Policy 进一步让 trainer 和 rollouter 同时跑。系统问题也随之变化:不再只是“怎么少搬数据”,而是“怎么在吞吐、新鲜度和 off-policy 修正之间取平衡”。

1. 入口:trainer 和 rollouter 同时运行

Fully async 的入口不是普通 RayPPOTrainer.fit()FullyAsyncTaskRunner会先创建 FullyAsyncTrainer,再创建 FullyAsyncRollouter,随后用 rollouter 给出的 max_queue_size创建 MessageQueue,并把同一个 MessageQueueClient注入两边。初始化完成后,它先做一次参数同步,再同时启动 rollouter.fit.remote()trainer.fit.remote()verl/experimental/fully_async_policy/fully_async_main.py:51-115176-209)。

下面这张图画的是运行入口。读图时注意:trainer 和 rollouter 不是一个函数调用栈里的上下游,而是两个 actor,各自长期运行,中间只通过 queue 和参数同步协议交互。

FullyAsyncTaskRunner 同时启动 trainer 和 rollouter

这解释了 fully async 的第一层工程含义:它不是把同步 step 里的某个阶段改快,而是改变阶段之间的依赖形态。同步 PPO 里 trainer 调用 rollout 并等待返回;这里 rollouter 独立生产样本,trainer 独立消费样本,失败或结束由 ray.wait()监控,最后清空 message queue。

2. MessageQueue 把 batch 返回改成样本消费

MessageQueue是一个 Ray actor,内部用 deque(maxlen=max_queue_size)存样本。put_sample()在队列满时会丢弃最旧样本并计数,随后追加新样本;get_sample()会等到队列非空或关闭,然后弹出最旧样本;get_statistics()暴露 queue size、produced、consumed、dropped 和 max size(verl/experimental/fully_async_policy/message_queue.py:26-119)。

trainer 侧不再等待某个 rollout 调用返回完整 DataProto。FullyAsyncTrainer.__init__()里把 required_samples定义成 actor.ppo_mini_batch_size * require_batches_get_samples_from_queue()会循环 message_queue_client.get_sample(),直到拿够样本,再用 assemble_batch_from_rollout_samples()拼成训练 batch,并记录等待时间(verl/experimental/fully_async_policy/fully_async_trainer.py:143-154273-330)。

下面这张图展示 producer/consumer 合同。读图时注意,队列里不是“一个 step 的完整 batch”,而是一条条 RolloutSample;trainer 每次消费多少,由 required_samples决定。

MessageQueue 把 rollout 输出变成样本消费

这个设计把吞吐问题拆成两类可观测状态:如果 trainer 等队列,fully_async/total_wait_time会上升;如果 rollouter 生产过快,queue size 和 dropped samples 会上升。第 26 篇里的 gen瓶颈在这里变成 producer/consumer 速率匹配问题,而不是单个阶段窗口。

3. staleness 是 fully async 的核心状态

rollouter 侧也不是无限生成。初始化时,它要求 train_batch_size == 0gen_batch_size == 1,并读取 async_training.staleness_thresholdtrigger_parameter_sync_steprequire_batchesset_max_required_samples()会把 required_samples * (staleness_threshold + 1) * trigger_parameter_sync_step作为 max_required_samples,同时把 max_queue_size设为这个值,并把并发样本数限制在 rollout replicas 数量和该上限之间(verl/experimental/fully_async_policy/fully_async_rollouter.py:392-419482-542)。

下面这张图要看的重点是 staleness 的生命周期:rollouter 每提交一个样本处理就增加 staleness 计数;如果队列满或 staleness 达到上限,就暂停继续提交;trainer 参数同步后再通知 rollouter reset。

Rollouter 用 staleness 控制继续生产

源码路径和这张图一致。_feed_samples()持续从 dataloader 取单样本,封装成 RolloutSample放进 pending_queue_processor_worker()从 pending queue 取样本,增加 staleness_samples,并在 active tasks 不超过 max_concurrent_samples时创建 _process_single_sample_streaming()任务;生成完成后把序列化后的样本写入 MessageQueueverl/experimental/fully_async_policy/fully_async_rollouter.py:805-940)。

暂停条件集中在 _should_pause_generation():如果 message queue 已满,或 staleness_samples >= max_required_samples,rollouter 就暂停。trainer 侧 _fit_update_weights()只在 local_trigger_step == 1时同步参数,同步后调用 rollouter.reset_staleness();rollouter 会把 staleness 重新设为当前 active tasks 加 queue size,因为这些样本仍可能来自旧参数版本(verl/experimental/fully_async_policy/fully_async_rollouter.py:562-5901067-1089verl/experimental/fully_async_policy/fully_async_trainer.py:487-525)。

所以 staleness 不是一个日志装饰项,而是 fully async 能不能成立的系统阀门。阈值太紧,rollouter 容易停下来,吞吐重叠变少;阈值太松,trainer 会消费更多旧版本样本,off-policy 压力上升。

4. partial rollout 让长尾不必全量作废

Fully async 还有一个容易忽略的细节:生成过程可能跨参数版本。FullyAsyncLLMServerClient.generate()会在 rollout 被 abort 时,根据 async_training.partial_rollout决定是否带着已有 output tokens 继续生成;最终输出会写入 global_stepsmin_global_stepsmax_global_stepsverl/experimental/fully_async_policy/fully_async_rollouter.py:51-150)。

下面这张图解释 partial rollout 的作用。读图时注意,它不是为了让旧样本消失,而是为了保留已经生成的前缀,并把跨版本跨度记录下来。

partial rollout 记录跨参数版本的轨迹跨度

拼 batch 时,assemble_batch_from_rollout_samples()会读取 min_global_stepsmax_global_steps,计算 partial ratio、max partial span、param version diversity,并把 trajectory_param_versions写入 meta_info。trainer 的 _collect_metrics_from_samples()会根据 current_param_version - trajectory_param_versions累计 stale trajectory 数(verl/experimental/fully_async_policy/detach_utils.py:84-176verl/experimental/fully_async_policy/fully_async_trainer.py:741-757)。

这就是 partial rollout 和 staleness 的互补关系:staleness 管“最多允许多少旧样本在系统里”,partial rollout 管“一个长生成被中断后如何继续,以及如何记录跨版本事实”。前者是流控,后者是样本语义记录。

5. off-policy correction 是速度账的算法侧收口

一旦 trainer 消费旧版本样本,系统就要面对 off-policy。verl 把这件事拆成两种模式:bypass_mode=True时直接把 rollout_log_probs作为 old_log_probs,跳过额外 old logprob forward;非 bypass 模式下,会重新计算 old logprob,并在 advantage 前通过 rollout correction helper 计算 IS weights、rejection mask 和诊断指标(verl/trainer/ppo/rollout_corr_helper.py:1102-1137verl/experimental/separation/ray_trainer.py:493-510560-599)。

下面这张图把算法侧代价放回系统三角形。读图时注意:速度来自异步队列,风险来自样本旧,修正来自 logprob 差异、IS/RS 和 off-policy metrics。

off-policy correction 收口 fully async 的算法风险

rollout_corr_helper.py的模块说明直接把 off-policy 来源列成三类:rollout 和 training 实现的 policy mismatch、旧 checkpoint 样本带来的 model update staleness、以及一般分布漂移。核心函数会根据 old_log_probs - rollout_log_probs计算 importance sampling weights,可选 rejection sampling,并输出 KL、PPL、log PPL diff、chi-square、ESS 等指标;compute_rollout_correction_and_add_to_batch()会更新 response_mask,并在启用 IS 时把 rollout_is_weights加回 batch(verl/trainer/ppo/rollout_corr_helper.py:14-61520-655779-890897-1060)。

这部分不能被理解成“有 correction 就可以无限异步”。correction 是把已经产生的 off-policy 差异显式度量和部分校正;staleness threshold、parameter sync cadence、queue size 和 partial rollout 才是在系统侧限制差异继续扩大的手段。

6. 读 fully async 指标时看这几组账

Fully async 的性能诊断不能只看 step time。它至少有四组账:

观察项

更可能说明什么

对应源码状态

fully_async/total_wait_time高

trainer 等不到足够样本

_get_samples_from_queue()等待 queue

queue size 高、dropped samples 高

rollouter 生产快于 trainer 消费,或 trainer 阶段太慢

MessageQueue.get_statistics()

count/staleness_samples接近上限

旧版本样本堆积,rollouter 可能暂停

_should_pause_generation()

partial ratio、max partial span 高

较多轨迹跨参数版本

assemble_batch_from_rollout_samples()

stale trajectory processed 高

trainer 正在消费旧版本轨迹

_collect_metrics_from_samples()

rollout correction KL、PPL、chi-square、ESS 异常

rollout policy 和 training policy 差距变大

rollout_corr_helper.py

timing_s/param_sync高

参数同步本身成为新瓶颈

_fit_update_weights()

这张表的保守用法是:先判断 trainer 是等样本、rollouter 是被 staleness 暂停,还是参数同步慢;再看 off-policy 指标是否说明异步带来的旧样本已经影响训练分布。不要只用吞吐上升判断 fully async 成功,也不要只用某个 KL 指标判断系统失败。

小结:fully async 把吞吐优化变成控制问题

放回系列地图,第 27 篇把 data movement 从 GPU 之外拎出来,第 28 篇把 controller 回流改成队列化数据系统,第 29 篇则把 trainer/rollouter 的 step 级同步进一步拆开。Fully Async Policy 的价值是让 rollout 生成、训练 update 和参数同步尽量重叠;它的约束是样本新鲜度不能再被默认保证。

所以本文的结论是:fully async 不是单纯的性能开关,而是一套控制系统。MessageQueue管 producer/consumer,staleness threshold 管旧样本上限,partial rollout 管跨版本轨迹记录,parameter sync 管版本推进,rollout correction 管算法侧诊断和修正。下一篇第 30 篇进入 671B/MoE 级别的 scale-out 场景时,这个三角关系会继续放大:吞吐、新鲜度和并行/显存/通信约束会一起成为系统设计的主轴。

本文源码索引

  • verl/experimental/fully_async_policy/fully_async_main.py:51-115176-209:FullyAsyncTaskRunner 如何创建 trainer、rollouter、MessageQueue,并同时运行两侧 actor。
  • verl/experimental/fully_async_policy/message_queue.py:26-119:MessageQueue 的 put/get/statistics 行为,以及队列满时丢弃最旧样本。
  • verl/experimental/fully_async_policy/fully_async_trainer.py:143-154273-330:trainer 如何定义 required_samples,并从 queue 凑 batch。
  • verl/experimental/fully_async_policy/fully_async_trainer.py:454-525:trainer 如何拉样本、维护 local trigger step、同步参数并 reset staleness。
  • verl/experimental/fully_async_policy/fully_async_rollouter.py:392-542:rollouter 的 fully async 配置约束、staleness threshold、max queue size 和 max concurrent samples。
  • verl/experimental/fully_async_policy/fully_async_rollouter.py:805-9401067-1089:rollouter 如何持续 feed/process 样本,以及何时因 queue 或 staleness 暂停。
  • verl/experimental/fully_async_policy/fully_async_rollouter.py:51-150562-590:partial rollout 如何记录参数版本跨度,以及参数同步后如何 reset staleness。
  • verl/experimental/fully_async_policy/detach_utils.py:84-176179-240:RolloutSample 拼 batch、partial 统计和 fully async metrics 聚合。
  • verl/experimental/separation/ray_trainer.py:493-510560-599:bypass/decoupled logprob 模式,以及 rollout correction 在 advantage 前的接入点。
  • verl/trainer/ppo/rollout_corr_helper.py:14-61520-655779-890897-10601102-1137:rollout correction 的 off-policy 来源、IS/RS、诊断指标和 bypass mode 行为。
本文参与 腾讯云自媒体同步曝光计划,分享自微信公众号。
原始发表:2026-07-24,如有侵权请联系 cloudcommunity@tencent.com 删除
目录
  • 1. 入口:trainer 和 rollouter 同时运行
  • 2. MessageQueue 把 batch 返回改成样本消费
  • 3. staleness 是 fully async 的核心状态
  • 4. partial rollout 让长尾不必全量作废
  • 5. off-policy correction 是速度账的算法侧收口
  • 6. 读 fully async 指标时看这几组账
  • 小结:fully async 把吞吐优化变成控制问题
  • 本文源码索引
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档