跳转至

env_loop.py — EnvLoop 管理模型与向量化环境之间的交互

文件路径: verl/experimental/vla/env_loop.py 模块路径: verl.experimental.vla.env_loop

文件概述

EnvLoop 管理模型与向量化环境之间的交互。它实现了流水线化(pipelined)的执行方式,通过将环境分成多个"阶段"(stage),让模型推理和环境仿真可以部分重叠执行,从而提高 GPU 利用率。

核心设计:流水线化交互

为什么需要流水线?

在机器人仿真中,环境步进(step)通常在 CPU 上运行(物理模拟),而模型推理在 GPU 上运行。如果串行执行:

[模型推理] -> [环境步进] -> [模型推理] -> [环境步进] -> ...
  GPU空闲      CPU空闲       GPU空闲       CPU空闲

通过流水线化,可以让两者重叠:

Stage A: [模型推理] -> [环境步进A] -> [模型推理] -> [环境步进A]
Stage B:              [模型推理] -> [环境步进B] -> [模型推理]
         GPU满载       GPU满载       GPU满载       GPU满载

初始化

class EnvLoop:
    def __init__(self, env_wg, rollout_wg, config):
        self.env_wg = env_wg          # 环境 worker group
        self.rollout_wg = rollout_wg  # 模型推理 worker group
        self.max_interactions = config.env.train.max_episode_steps // \
                               config.env.actor.model.num_action_chunks
        self.stage_num = config.env.rollout.pipeline_stage_num
        self.total_envs = self.env_wg.world_size * self.num_envs_per_worker
        self.envs_per_stage = self.total_envs // self.stage_num

核心方法 rollout

async def rollout(self):
    """执行一轮完整的环境交互"""
    # 为每个阶段创建异步任务
    tasks = []
    for stage_id in range(self.stage_num):
        task = asyncio.create_task(
            self._stage_loop(stage_id)
        )
        tasks.append(task)

    # 并行执行所有阶段
    results = await asyncio.gather(*tasks)

    # 合并所有阶段的轨迹
    trajectories = self._collate_trajectories(results)
    return trajectories

单阶段循环

async def _stage_loop(self, stage_id):
    """单个阶段的交互循环"""
    trajectory = []
    for step in range(self.max_interactions):
        # 1. 获取当前观测
        obs = self._restructure_obs_data(stage_id)

        # 2. 模型推理(GPU)
        actions = await self.rollout_wg.generate_sequences(obs)

        # 3. 环境步进(CPU,带 stage_id)
        env_output = await self.env_wg.env_interact_step(
            actions, stage_id=stage_id
        )

        # 4. 记录转移
        trajectory.append((obs, actions, env_output))

    return trajectory

数据重组

def _restructure_obs_data(self, stage_id):
    """将环境输出重组为模型输入格式

    环境返回: {full_image, wrist_image, state, task_descriptions}
    模型需要: DataProto 格式的输入
    """
    ...

def _collate_trajectories(self, results):
    """合并多个阶段的轨迹

    将各阶段的独立轨迹拼接成完整的训练数据。
    """
    ...

流水线示意图

时间 -->
Stage 0: [推理0] [环境0] [推理0] [环境0] ...
Stage 1:          [推理1] [环境1] [推理1] [环境1] ...
Stage 2:                   [推理2] [环境2] [推理2] ...

GPU:     [推理0] [推理1] [推理2] [推理0] [推理1] ...
         ^--- GPU 始终保持忙碌

核心类/函数列表

名称 类型 说明
EnvLoop 类 流水线化环境交互管理器
rollout 方法 执行完整的交互循环
_stage_loop 方法 单阶段交互循环
_restructure_obs_data 方法 观测数据格式转换
_collate_trajectories 方法 合并多阶段轨迹

与其他模块的关系

  • 被 RobRayPPOTrainer(rob_ray_trainer.py)调用
  • 使用 EnvWorker(workers/env/env_worker.py)的 env_interact_step
  • 使用 Rollout Worker 的 generate_sequences 做模型推理

小结

EnvLoop 是 VLA 训练中提高效率的关键组件。通过流水线化,它让 GPU(模型推理)和 CPU(环境仿真)尽可能并行工作,避免资源空闲。stage_num 参数控制流水线的深度,通常设为 2-4。