verl SkipManager 完全指南:在 RL 训练流水线中按步跳过 rollout 的缓存与回放机制
verl SkipManager 完全指南在 RL 训练流水线中按步跳过 rollout 的缓存与回放机制【免费下载链接】verlverl/HybridFlow: A Flexible and Efficient RL Post-Training Framework项目地址: https://gitcode.com/GitHub_Trending/ve/verlverlHybridFlow的 SkipManager 提供了一套通用化的“跳过机制”允许开发者针对训练流水线中的指定步骤step跳过代价高昂的 rollout 生成阶段从而节省时间、显存等资源并显著加快调试与迭代速度。本文以 docs/advance/skip_manager.rst 为核心骨架结合 verl/utils/skip 下的源码实现系统讲解 skip 配置、cache/repeat两种动作语义、三种内置角色rollout、rollout_tq、async_rollout的接入方式、V1 双阶段装饰器原理以及如何扩展自定义跳过模块。1. 概述SkipManager 是什么SkipManagerverl.utils.skip.SkipManager是一个通用化的“跳过框架”用于在 verl 训练流程中跳过选定步骤。通过在配置好的步骤上绕过昂贵阶段帮助节省时间、显存或其他资源并在调试与实验过程中提升开发者迭代速度。跳过的行为统一收敛在 Hydra 顶层键skip之下。模块按**角色role**注册例如rollout、rollout_tq、async_rollout并通过SkipManager.annotate(role...)装饰器挂接到目标函数上V1 两阶段路径使用SkipManager.annotate_tq(role..., phase...)。每个角色在配置中声明哪些整数 step 允许触发跳过逻辑。当前仅实现了与 rollout 相关的角色同一机制可扩展到流水线的其他阶段见第 6 节。1.1 典型使用场景SkipManager 面向“重复完整训练代价高昂”的开发工作流加速迭代在选定 step 上跳过重阶段如生成同时继续训练其余流水线。确定性回放Deterministic replay缓存并重载中间结果以便在特定 step 上复现此前某次运行。资源节省在二分定位 bug 或调试下游逻辑时避免重复计算或长时间持有大张量。内置的rollout/rollout_tq/async_rollout模块将该机制应用于序列生成其他角色可以按同样的模式新增。1.2 当前支持的入口点训练入口Skip 角色 / 配置状态main_ppo.py且trainer.use_v1FalseRayPPOTrainerskip.rolloutSupportedmain_ppo.py且trainer.use_v1TrueV1PPOTrainer TransferQueueskip.rollout_tqSupported见第 4 节fully_async_mainFullyAsyncRollouterskip.async_rolloutSupported2. 配置说明skip.rollout/skip.rollout_tq/skip.async_rollout三种角色使用同一套 Hydra 字段集对应的数据类为RolloutSkipConfig/RolloutTqSkipConfig/AsyncRolloutSkipConfig定义于 verl/utils/skip/config.py。默认值位于 verl/trainer/config/ppo_trainer.yaml 中各自的skip.*键下。2.1 参数一览参数类型说明enablebool该角色的总开关。默认False。dump_dirstr缓存分片shard的根目录支持~展开。默认~/.verl/rollout_dump。stepslist[int]允许触发跳过逻辑的 step 列表。不在此列表内的 step 上被装饰函数总是正常运行。actioncache|repeat跳过时采取的动作见下文。默认cache。steps的语义随角色不同而不同对skip.rollouttrainer 的global_steps通过SkipManager.set_step设置。对skip.rollout_tqtrainer 的global_steps通过SkipManager.set_step设置。对skip.async_rollout从sample_id解析出的feed-order 索引见第 5 节——不是trainer 的global_steps。2.2 action 语义cache若当前 step 存在有效缓存 dump则加载它并跳过生成否则运行生成并把结果写入该 step 目录。repeat若存在任何有效 dump则按下方算法从替代 step加载否则正常运行生成并照常 dump。[!NOTE] 当前配置中只校验cache与repeat两种 action。尽管 base_skip.py 中的SkipAction枚举还列出了RANDOM、EMPTY等值供未来模块使用。配置数据类在__post_init__中会进行严格校验见 config.pyenable必须是 bool、dump_dir必须是 str、steps必须全为 int且action必须属于{cache, repeat}否则启动即报错。2.3 repeat 的替代 step 选择算法rollout/async_rolloutRolloutSkip._find_latest_step当actionrepeat且当前 step 目录缺失或不完整时若当前step 的目录有效直接使用当前 step。否则使用严格小于当前 step 的最大可用 step。否则使用严格大于当前 step 的最小可用 step。若不存在任何有效 dump则跳过不适用被包装函数正常运行之后可能 dump。源码实现位于 rollout_skip.py_get_available_steps扫描项目 dump 目录下所有可解析为 int 且通过_check_valid_step_path校验的子目录并排序_find_latest_step依序选择精确命中、左侧最近、右侧最近最后返回-1。[!WARNING]repeat不保证缓存批次与当前 prompt 或 trainer step 对齐——仅用于调试与迭代需要 step 对齐回放时请优先使用cache。rollout_tqRolloutTqSkip._resolve_load_step_v1与上述回退策略一致但检查的是 V1 格式缓存文件tq_batch.pt见 2.5 节磁盘布局而非gen_batch.dp。其实现见 rollout_skip.py。2.4 Hydra CLI 配置示例Colocated PPOskip.rolloutskip.rollout.enableTrue skip.rollout.dump_dir/path/to/rollout_dump skip.rollout.steps[1,2,3,10] skip.rollout.actioncache基于 TransferQueue 的 V1 trainerskip.rollout_tqskip.rollout_tq.enableTrue skip.rollout_tq.dump_dir/path/to/rollout_dump skip.rollout_tq.steps[1,3,5] skip.rollout_tq.actioncacheFully asyncskip.async_rolloutskip.async_rollout.enableTrue skip.async_rollout.dump_dir/path/to/rollout_dump skip.async_rollout.steps[1,2,3,4,5] skip.async_rollout.actioncache若只想在bash中传长 step 列表静态 YAML 中不可写可使用 shell 展开skip.async_rollout.steps[$(seq -s, 1 128)]2.5 磁盘目录布局所有角色共享同一套项目级目录结构仅每个 step 内的文件不同{dump_dir}/{experiment_name}_{project_name}/ └── GBS{gbs}_N{n}_in{prompt_len}_out{response_len}/ ├── {step}/ │ ├── gen_batch.dp # rollout / async_rollout │ ├── tq_batch.pt # rollout_tq │ └── meta.json └── ...experiment_name/project_name取自运行配置中的trainer.experiment_name与trainer.project_name。gbs、n、prompt_len、response_len分别来自data.gen_batch_size或训练 batch size、actor_rollout_ref.rollout.n、data.max_prompt_length、data.max_response_length。这些元信息在源码中由 RolloutSkip.__init__ 通过OmegaConf.select从全局配置中解析并由_get_project_dump_dir拼装目录名。[!CAUTION] colocatedmain_ppo的缓存通常GBS 较大与 fully async 流式缓存的缓存通常GBS1在元信息不一致时一般不可互换。rolloutgen_batch.dp与rollout_tqtq_batch.pt使用不同文件格式即使项目元信息一致也永远不可互换。角色缓存文件内容rollout / async_rolloutgen_batch.dp通过DataProto.save_to_disk保存的 DataProto——完整的generate_sequences输出prompts、responses、log_probs 等。rollout_tqtq_batch.pttorch.save载荷包含tensordict通过tq.kv_batch_get读取的全部轨迹字段、tags逐轨迹 tag 列表、keys轨迹级 TQ 键、global_stepsint。详见第 4 节。2.6 最小工作流cache 模式首次运行设置enableTrue、actioncachesteps列出关心的 step。dump_dir为空 → 正常运行生成并为每个 step 写入缓存文件 meta.json。第二次运行使用相同配置与兼容的 trainer 元信息 → 列表内的 step 从磁盘加载而非重新生成。部分缓存部分 step 目录缺失缺失的 step 会在下次运行重新生成其他 step 若存在缓存仍会加载。2.7 与旧版 RolloutSkip 的关系如果skip.rollout.enable与旧版actor_rollout_ref.rollout.skip.enable同时为 trueSkipManager 会发出DeprecationWarning并强制将旧版标志置为False保证只有一个机制在运行。3. Rollout 角色快速上手rolloutrole当使用main_ppo.py且trainer.use_v1FalseRayPPOTrainer走标准AgentLoopManager.generate_sequences路径时使用skip.rollout。配置字段与cache/repeat语义见第 2 节。接线Wiringray_trainer.py 中RayPPOTrainer.fit()调用SkipManager.init(self.config)并在每个训练 step 调用SkipManager.set_step(self.global_steps)见 L1455。agent_loop.py 中AgentLoopManager.generate_sequences被SkipManager.annotate(rolerollout)装饰。被装饰函数是单一入口点它接收 prompts驱动全 batch 生成chunk 分发、concat、计时返回完整 DataProto。跳过逻辑把这个单元整体包装——cache 命中时直接返回缓存 DataProto不调用 LLMcache 未命中时运行生成并 dump 结果。skip_manager.py中的annotate装饰器会自动识别同步与协程函数inspect.iscoroutinefunction分支并带有一个_should_bypass_for_validation校验旁路当批次meta_info或non_tensor_data标记validateTrue验证/评测路径时始终直通原函数不会被跳过逻辑拦截。4. Trainer V1 快速上手rollout_tqrole当使用main_ppo.py且trainer.use_v1TrueV1PPOTrainer TransferQueue时使用skip.rollout_tq。该配置覆盖全部 V1 trainer 模式sync、colocate_async、separate_async。[!IMPORTANT] 与rollout第 3 节不同V1 TransferQueue 路径没有单一的generate_sequences入口点。Rollout 被拆分到训练 step 中不同时机执行的两个方法提交阶段Submit phase——PPOTrainer._add_batch_to_generate从 dataloader 采样 batch、分配 uid、将 prompts 派发给AgentLoopManager进行生成。采样阶段Sample phase——ReplayBuffer.sample等待轨迹完成然后从 TQ 收集为KVBatchMeta。单一装饰器无法同时覆盖两者因为 cache 命中的短路必须发生在提交阶段内部——在 uid 生成之后uid 用于键映射、真实 rollout 提交之前这正是想跳过的部分。解决方案是SkipManager.annotate_tq一个按phasesubmit或phasesample选择的两阶段装饰器。4.1 接线trainer_base.py 中PPOTrainer.fit()调用SkipManager.init(self.config)并在每个训练 step 调用SkipManager.set_step(self.global_steps)见 L433、L499。trainer_base.py 中PPOTrainer._add_batch_to_generate被SkipManager.annotate_tq(rolerollout_tq, phasesubmit)装饰。replay_buffer.py以及 L581中ReplayBuffer.sample被SkipManager.annotate_tq(rolerollout_tq, phasesample)装饰。4.2 方法拆分为了让提交阶段装饰器有一个干净的拦截窗口_add_batch_to_generate被拆成两个子方法SkipManager.annotate_tq(rolerollout_tq, phasesubmit) def _add_batch_to_generate(self): batch self._next_train_batch() # dataloader uid assignment self._submit_batch_to_rollout(batch) # tag registration generate_sequences def _next_train_batch(self): Advance the dataloader and return a batch with fresh uids. ... def _submit_batch_to_rollout(self, batch): Register prompt tags in TransferQueue and dispatch to AgentLoopManager. ...当 skip禁用或当前 step 不在skip.rollout_tq.steps中时装饰器直通函数体正常执行_next_train_batch再_submit_batch_to_rollout。当 skip启用且 step 可跳过时装饰器接管函数体不执行。装饰器自己调用_next_train_batch即使在 cache 命中 step 上也保持 dataloader 对齐然后按缓存可用性分支。[!NOTE] 若未来改动在_add_batch_to_generate内部_next_train_batch与_submit_batch_to_rollout之间新增逻辑该逻辑也必须同步反映到装饰器的 submit 阶段分支中因为 skip 激活时装饰器绕过了函数体。4.3 两阶段流程阶段 1 —— cache 未命中首次运行磁盘无tq_batch.ptstep() | - _add_batch_to_generate() [submit decorator: skip enabled, cache-miss] | - _next_train_batch() - batch with fresh uids | - maybe_load_and_inject() - False (no cache) | - _submit_batch_to_rollout() - real LLM generation dispatched to TQ | - replay_buffer.sample() [sample decorator: skip enabled] | - (original sample runs) - waits for trajectories, returns KVBatchMeta | - should_save() - True (no cache, partitiontrain) | - prepare_data() - kv_batch_get all fields - torch.save - tq_batch.pt | - (downstream: reward, advantage, actor/critic update ...)阶段 2 —— cache 命中再次运行tq_batch.pt存在step() | - _add_batch_to_generate() [submit decorator: skip enabled, cache-hit] | - _next_train_batch() - batch with fresh uids | - maybe_load_and_inject() - True - load_dump_data() | | - torch.load(tq_batch.pt) - old keys, tags, tensordict | | - group old trajectories by uid prefix | | - map new uids - cached groups (modulo cycling) | | - index_select_tensor_dict - select/repeat trajectory rows | | - kv_batch_put trajectories with new keys updated tags | | - kv_batch_put prompt-level keys with statusfinished | - return (skip _submit_batch_to_rollout -- no real LLM call) | - replay_buffer.sample() [sample decorator: skip enabled] | - (original sample runs) - finds finished prompts immediately, returns KVBatchMeta | - should_save() - False (cache exists) | - (no prepare_data call) | - (downstream: reward, advantage, actor/critic update ...)4.4 磁盘格式tq_batch.ptprepare_data实现见 rollout_skip.py保存一个含四个键的torch.save载荷键内容tensordict一个TensorDict包含通过tq.kv_batch_get(keysbatch.keys)从 TQ 读取的全部轨迹字段无select_fields过滤。包括prompts、responses、response_mask、数据集字段messages、datasource等以及 agent loop 写入的任何字段。其中prompts/responses是NestedTensorjagged锯齿形其余为常规张量。tagslist[dict]——逐轨迹 tag与keys中的条目一一对应且顺序一致。每个 tag 携带global_steps、status、seq_len等。keyslist[str]——{uid}_{session_id}_{index}格式的轨迹级 TQ 键。uid 前缀是 UUID4不含下划线因此key.split(_)[0]可以还原出父级 prompt uid。global_stepsint——dump 创建时的 step用于健全性检查。meta.json记录{global_steps: int, num_trajectories: int}被_check_valid_v1_step_path用于完整性校验rollout_skip.py。4.5 为什么kv_batch_get不带select_fieldsReplayBuffer.sample返回的KVBatchMeta只持有键/tag 引用——轨迹数据实体存放在 TQ 的存储单元中。TQ 会在每个 step 结束时清理键PPOTrainer.fit中的tq.kv_clear因此prepare_data必须在清理之前读取完整数据实体。不带select_fields读取可确保所有下游字段reward、log-prob、masks都被捕获以进行忠实的回放。4.6 Cache 命中注入键重映射load_dump_data实现见 rollout_skip.py不能直接复用旧键——_next_train_batch刚生成了新的 uid。它按如下步骤把缓存轨迹重映射到新 uid 上按 uid 分组解析old_keys按key.split(_)[0]将轨迹索引分组。每组对应一个 prompt GRPO 组通常为n条轨迹但若部分 session 失败可能更少。这里使用 dict 分组因此不依赖键按 uid 前缀排序。将新 uid 映射到组对第prompt_idx个新 uid选择groups[prompt_idx % num_cached_groups]。若缓存组少于新 prompt 数组会按模循环若更多只使用前num_prompts个组。每个 prompt 填充n条轨迹对[0, n)中的每个session_id取group[session_id % len(group)]。若组内条目少于n部分失败轨迹在组内循环以填满槽位。构造新键{new_uid}_{session_id}_0——匹配标准 TQ 键格式。更新 tags每条轨迹 tag 的global_steps/min_global_steps/max_global_steps被覆写为当前step使ReplayBuffer的新鲜度检查_drop_max_off_policy_samples不会把它们当作 off-policy 丢弃。同时从轨迹 tag 中移除is_prompt标志该标志只属于 prompt 级键。选择张量行index_select_tensor_dict(data, traj_indices)从缓存tensordict中选取并可能复制行。该函数同时处理常规张量与NestedTensorunbind - select -nested_tensor_from_tensor_list重建因为NestedTensor不支持在 dim0 上直接索引/切片。两次kv_batch_put轨迹数据keysnew_keys、fieldsnew_fields、tagsnew_tags——将真实轨迹内容写入 TQ 存储。Prompt 级标记keysnew_prompt_uids、tags[{is_prompt: True, status: finished, ...}]——将每个新 prompt 标记为已完成使ReplayBuffer.sample的_has_enough_samples无需轮询即可立即通过。注入完成后TQ 中的结构与一次正常完成的 rollout 完全相同每个 prompt 对应n个轨迹键 一个statusfinished的 prompt 级键。下一次调用replay_buffer.sample即可拾取它们返回数据为注入缓存内容的KVBatchMeta。[!TIP] 在parameter_sync_step 1separate async时一个 global step 内会多次调用sample。每个 mini-batch 会保存到独立的内部子目录{step}/{inner_idx}/加载时合并全部子目录parameter_sync_step 1sync / colocate async时目录结构为{step}/{0}/tq_batch.pt见 RolloutTqSkip 类注释。5. Fully async 快速上手async_rolloutrole在 fully_async.md 描述的架构中Trainer 与 Rollouter 运行在不同进程。Rollout 生成发生在 Rollouter 上采用流式单样本分发。启动fully_async_main时应使用skip.async_rollout而不是skip.rollout。共享的 Hydra 字段与磁盘布局见第 2 节。[!IMPORTANT] 在async_rollout中step不是trainer 的时间线。它只是 Rollouter 上的prompt 请求/喂入顺序feed order即FullyAsyncRollouter入队下一个 prompt 时sample_{epoch}_{index}中的单调索引。并发 rollout 下完成顺序可能不同于喂入顺序配置skip.async_rollout.steps时不要把这些索引当作 trainerglobal_steps或参数同步边界。5.1 从sample_id提取 step 键每个被喂入的样本携带形如sample_{epoch}_{index}的 id例如sample_0_42。与skip.async_rollout.steps匹配、并用于磁盘目录的整数是最后一段——入队时 Rollouter 的喂入顺序索引。源码中parse_async_rollout_sample_steprollout_skip.py支持sample_epoch_feed_index及带uid_前缀两种格式并严格校验三段式结构。5.2 接线fully_async_rollouter.py 中FullyAsyncRollouter在 Rollouter 进程内调用SkipManager.init(self.config)。fully_async_rollouter.py 中FullyAsyncAgentLoopManager.generate_sequences_single被SkipManager.annotate(roleasync_rollout)装饰并通过sample_id做在线 step 解析。AsyncRolloutSkip.extract_steprollout_skip.py从prompts.non_tensor_batch[uid]读取 uid 并解析出 feed index使每个并发样本独立解析 step。6. 设计与实现6.1 SkipManager APISkipManagerverl/utils/skip/skip_manager.py是一个类级注册表class-level registry核心 API 如下init(config)将config.skip解析为SkipManagerConfig为每个已注册角色实例化一个 skip 模块并存入SkipManager.skip_instances。实现上通过omega_conf_to_dataclass(config.skip, dataclass_typeSkipManagerConfig)完成 Hydra 配置到数据类的转换再遍历SKIP_REGISTRY实例化。set_step(step: int)为support_online_step False的角色设置SkipManager.stepmain_ppo与 V1PPOTrainer中使用 trainerglobal_steps。annotate(role, **kwargs)同步或异步函数的装饰器工厂rollout与async_rollout使用。annotate_tq(role, phase)V1 TransferQueue 路径的两阶段装饰器工厂rollout_tq使用。见第 4 节。SkipManager的类属性在init之前有默认值step-1、skip_instances{}因此仅导入装饰器尚未初始化的代码路径如测试也能工作——此时装饰器是 no-op。6.2 装饰器流程annotaterollout/async_rolloutcall decorated function │ ▼ skip disabled or role missing? ──yes──► run original function │no ▼ resolve step (set_step vs extract_step) │ ▼ step ∉ config.steps? ──yes──► run original function │no ▼ meet_precondition (cache/repeat)? ──yes──► warp_function (load cache) │no ▼ run original function → prepare_data (dump)该流程对应 skip_manager.py 中annotate的同步/异步两个 wrapper依次检查校验旁路、角色是否启用、step 解析与命中判断再进入meet_precondition/warp_function/prepare_data分支。6.3 BaseSkip 接口每个 skip 模块继承BaseSkipverl/utils/skip/base_skip.py并通过register_skip(role_name)注册进SKIP_REGISTRY。support_actions该模块允许的SkipAction值。RolloutSkip/RolloutTqSkip/AsyncRolloutSkip均声明为[CACHE, REPEAT]BaseSkip.__init__会在 action 不受支持时抛出ValueError。support_online_step为True时每次调用使用extract_step而非SkipManager.step。实例方法is_enabled、meet_precondition、warp_function、prepare_data以及support_online_stepTrue时必需的extract_step。RolloutSkip/RolloutTqSkip/AsyncRolloutSkip均在 verl/utils/skip/rollout_skip.py 中为三种角色实现了生成缓存。RolloutTqSkip继承RolloutSkip并新增 V1 专属方法should_save、maybe_load_and_inject、load_dump_data、has_v1_cache、_resolve_load_step_v1。6.4 被拦截的函数角色被装饰函数定义位置step 来源rolloutAgentLoopManager.generate_sequencesverl/experimental/agent_loop/agent_loop.pySkipManager.set_step- trainerglobal_stepsrollout_tqPPOTrainer._add_batch_to_generatephasesubmitverl/trainer/ppo/v1/trainer_base.pySkipManager.set_step- trainerglobal_stepsrollout_tqReplayBuffer.samplephasesampleverl/trainer/ppo/v1/replay_buffer.pySkipManager.set_step- trainerglobal_stepsasync_rolloutFullyAsyncAgentLoopManager.generate_sequences_singleverl/experimental/fully_async_policy/fully_async_rollouter.pyextract_step→sample_id后缀 →prompt 喂入顺序rollout将完整 batch 的 Agent Loop RPCchunk 分发、concat、计时作为一个跳过单元整体包装。rollout_tq通过单个annotate_tq装饰器按phase选择包装两个方法submit 阶段在真实生成之前拦截sample 阶段在采样之后拦截以持久化结果。async_rollout包装单条流式样本的generate_sequences_single(self, prompts, sample_id)使并发样本独立解析 step。6.5 Step 解析set_stepvssupport_online_step共享SkipManager.step每进程一个类级槽位适用于顺序训练循环main_pporollout 前调用set_step(global_steps)。在线 stepAsyncRolloutSkip设置support_online_step True每次调用解析sample_id使在途的异步样本不共享单一计数器。对repeatRolloutSkip在每次meet_precondition与warp_function调用时重新计算_find_latest_stepskip 实例上不持有可变 step 字段。6.6 扩展自定义 skip 模块将同一机制扩展到其他流水线阶段的步骤从 verl/utils/skip/base_skip.py 继承BaseSkip。用register_skip(your_role_name)装饰该类。在SkipManagerConfig中新增对应字段参照 config.py 中rollout/async_rollout/rollout_tq三个子配置的写法。挂接SkipManager.annotate(roleyour_role_name)。对并发流水线优先设置support_online_step True并通过调用参数传递 step 标识。对拆分式架构路径如 V1 的 submit/sample 分离使用SkipManager.annotate_tq(role..., phase...)并将目标方法拆开使装饰器能在两个子步骤之间拦截。7. 小结SkipManager 通过“角色注册 装饰器包装 磁盘缓存”的三层设计把“跳过某一步的 rollout 生成”变成了纯配置驱动的能力cache提供 step 对齐的确定性回放repeat提供迭代调试时的替代批次复用rollout、rollout_tq、async_rollout三种角色分别覆盖 colocated PPO、V1 TransferQueue 与 fully async 三大训练路径。其扩展点BaseSkip子类 register_skipSkipManagerConfig字段为未来跳过 reward、advantage 等其他流水线阶段预留了清晰的实现范式。【免费下载链接】verlverl/HybridFlow: A Flexible and Efficient RL Post-Training Framework项目地址: https://gitcode.com/GitHub_Trending/ve/verl创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考