Ray RLlib EnvRunner API 深度解析:分布式环境采样、生命周期管理与容错机制 Ray RLlib EnvRunner API 深度解析分布式环境采样、生命周期管理与容错机制【免费下载链接】rayRay is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads.项目地址: https://gitcode.com/gh_mirrors/ra/rayEnvRunner 是 Ray RLlib 新 API 栈New API Stack中负责从环境中分布式采集训练数据的核心抽象它把环境创建、RLModule 推理、向量化采样、指标统计和容错恢复封装在一个可在 Ray Actor 中并行运行的单元内。本文以 env_runner.rst 参考文档为骨架结合 env_runner.py 等源码实现完整讲解 EnvRunner 的构造与设置、采样流程、指标与清理 API以及单智能体 / 多智能体两种内置子类的差异和基于StepFailedRecreateEnvError的子环境容错方案帮助你理解 RLlib 训练循环中数据从哪来、如何并行、如何容错这一关键环节。说明根据 new_api_stack.rst 的说明当前 Ray 版本默认使用 RLlib 新 API 栈EnvRunner正是新 API 栈中替代旧版RolloutWorker的数据采集组件。EnvRunner 的设计定位与核心职责在 RLlib 新 API 栈中EnvRunner位于 python/ray/rllib/env/env_runner.py是一个标记为PublicAPI(stabilityalpha)的抽象基类其核心功能可以概括为以下四点按配置构建组件通过向构造函数传入一个AlgorithmConfig对象完成初始化子类基于该配置构建可能是向量化的环境副本以及RLModule/Policy并用它们与环境交互来采集训练数据。对外提供采样入口使用者通过sample()方法从环境中收集训练数据。天然支持并行EnvRunner 可以通过ray.remote([resources])(EnvRunner)被封装为远程 Ray Actor随后用[ctor].remote(...)语法实例化 N 个 Actor实现多环境并行采样。可查询运行位置客户端可以获知每个 Actor 实际运行在哪个服务端 / 节点上用于亲和性与资源管理。从类继承关系看EnvRunner继承自FaultAwareApplypython/ray/rllib/utils/actor_manager.py并声明为abc.ABCMeta抽象类。它的__init__env_runner.py完成了以下基础工作复制配置self.config config.copy(copy_frozenFalse)每个 EnvRunner 持有自己的一份可变配置副本因此可以在运行期修改env_config、rl_module_spec等字段后重新构建环境或模块。记录 worker 身份通过worker_index和num_workers默认取config.num_env_runners标识当前实例在集群中的位置。初始化指标记录器self.metrics MetricsLogger(...)后续所有采样与回合统计都写入这里最终由get_metrics()汇总。确定性种子管理当config.seed非空时按公式seed (worker_index or 0) (1e6 * config.in_evaluation)计算每个 worker 的专属种子——评测evaluationworker 额外加1e6以避免与训练 worker 采样序列重叠然后通过update_global_seed_if_necessary统一设置 random / numpy / torch / tf 的全局种子env_runner.py。注册 Prometheus 计数器创建rllib_env_runner_num_try_env_step_counter尝试env.step()的次数与rllib_env_runner_num_env_steps_sampled_counter实际采样步数两个 Ray 指标计数器均带rllibtagenv_runner.py。类中还定义了模块级常量ENV_RESET_FAILURE、ENV_STEP_FAILURE和NUM_ENV_STEP_FAILURES_LIFETIME前者是_try_env_step在步骤失败时的返回哨兵值后者用于记录生命周期内累计的步进失败次数。构造与设置五个关键 API参考文档中Construction and setup一节列出了五个方法EnvRunner构造、make_env、make_module、get_spaces、assert_healthy。构造函数与make_env创建向量化环境make_env()的职责是创建 RL 环境并赋值给self.envenv_runner.py。它有两个重要的使用约定用户可以修改 EnvRunner 的配置例如改动self.config.env_config然后再次调用make_env()以新配置重建环境当既有环境发生故障时也应该调用它来关闭旧环境、重建新环境并继续采样。在SingleAgentEnvRunner中make_env()的实现single_agent_env_runner.py展示了完整的构建链路若已有环境则先尝试close()把self.config.env_config包装成EnvContext携带worker_index、num_workers、remote等上下文根据config.env的类型分三条路径注册环境字符串 ID 且已通过tune.registry注册过的环境、可调用对象callable、或者直接是 gymnasium 注册表中的字符串 ID最终通过gym.make_vec(env_name, num_envsself.config.num_envs_per_env_runner, vectorization_mode...)创建向量化环境并用DictInfoToList包装器将 info 字典转换为列表形式self.num_envs即向量化环境的并行子环境数量设置_needs_initial_reset True并在环境创建完成后触发on_environment_created回调。make_module构建推理用 RLModulemake_module()创建 EnvRunner 内部使用的RLModule并赋值给self.moduleenv_runner.py同样支持改配置后重建。单智能体实现通过self.config.get_rl_module_spec(envenv, spacesself.get_spaces(), inference_onlyTrue)获取模块规格注意inference_onlyTrueEnvRunner 中的模块只做前向推理不做训练再module_spec.build()并移动到设备上single_agent_env_runner.py。如果get_rl_module_spec抛NotImplementedError则self.module None——此时 EnvRunner 仍可配合random_actionsTrue做纯随机动作采样。get_spaces向算法暴露观测与动作空间get_spaces()是一个抽象方法返回一个将 ModuleID 映射到(观测空间, 动作空间)二元组的字典env_runner.py。单智能体实现返回三类键single_agent_env_runner.pyINPUT_ENV_SPACES向量化环境整体的(observation_space, action_space)INPUT_ENV_SINGLE_SPACES单个子环境的(single_observation_space, single_action_space)DEFAULT_MODULE_IDEnv-to-Module 连接器处理后的观测空间与单动作空间的组合。assert_healthy远程 Actor 的初始化探针assert_healthy()用于检查__init__()是否完整执行完毕特别适合 EnvRunner 以ray.remoteActor 方式运行时由属主进程确认 Actor 已正确初始化env_runner.py。单智能体实现断言self.env存在且module属性已挂载single_agent_env_runner.py不满足即抛AssertionError。采样与指标sample / get_metrics参考文档Sampling一节列出两个方法sample()与get_metrics()。sample()返回 Episode 数据sample(**kwargs)是抽象方法返回从该 EnvRunner 采集到的经验数据形式由配置决定。单智能体版本的完整签名与行为single_agent_env_runner.py如下def sample( self, *, num_timesteps: int None, # 本次最少采集多少时间步 num_episodes: int None, # 本次最少采集多少完整回合 explore: bool None, # None 时取 self.config.explore random_actions: bool False, force_reset: bool False, # 采样前强制重置所有向量化子环境 ) - List[SingleAgentEpisode]:几个关键行为细节num_timesteps与num_episodes二者只能指定其一同时传入会触发断言因为并行地从多个子环境采样实际返回的时间步 / 回合数会落在[请求值, 请求值 num_envs_per_env_runner - 1]区间内。两者都未指定且config.batch_mode truncate_episodes时默认采样config.get_rollout_fragment_length(worker_index) * num_envs步若batch_mode为complete_episodes则按train_batch_size驱动采样直到步数足够。exploreTrue时使用 RLModule 的forward_exploration()exploreFalse时使用forward_inference()随机动作模式下直接从env.action_space.sample()采样完全不经过模块。采样开始前若配置了全局EnvRunnerStateServerconfig.use_env_runner_state_serverTrue会先调用pull_if_newer(weights_seq_no)做一次只拉取更新版本的状态同步服务不可达时降级为继续使用当前权重并只告警一次single_agent_env_runner.py。采样结束后触发on_sample_end回调并把本次耗时记录进TIME_BETWEEN_SAMPLING。在_sample()内部循环single_agent_env_runner.py中每一轮迭代遵循典型的RL 数据流读取缓存的 Env-to-Module 连接器输出self._cached_to_module执行 RLModule 前向探索或推理经过 Module-to-Env 连接器得到最终动作调用_try_env_step(actions_for_env)推进环境若返回ENV_STEP_FAILURE哨兵值则重置所有环境并跳过本轮计数将(obs, action, reward, info, terminated, truncated)逐步写入对应的SingleAgentEpisode通过add_env_reset/add_env_step并在各阶段触发on_episode_created/on_episode_start/on_episode_step/on_episode_end回调回合结束时若episodes_to_numpyTrue默认则将 episode 转成 numpy 数组返回否则返回列表形式在num_timesteps模式下尚未结束的回合会通过episodes.cut(len_lookback_bufferconfig.episode_lookback_horizon)切成 chunk下一次sample()调用时从 chunk 续接既保证返回数据及时性又不丢失回合的连续性。每次sample()返回的是done_episodes_to_return ongoing_episodes_to_return即本次完成的回合与进行中的回合 chunk两部分single_agent_env_runner.py。配套的快速方法sample_get_state_and_metrics对于需要异步、低延迟的训练算法基类还提供了DeveloperAPI标注的sample_get_state_and_metrics()便捷方法env_runner.py它在同一次远程调用中完成sample()、取连接器状态get_state(not_componentsCOMPONENT_RL_MODULE)即排除模块权重和get_metrics()并把 episode 列表通过ray.put()存为ObjectRef返回——这样 episode 数据可以直接传给 Aggregator / Learner Actor而无需先经过主算法进程。get_metrics()回合级统计的汇总出口get_metrics()返回该 EnvRunner 到目前为止已采集完成回合的指标env_runner.py。单智能体实现single_agent_env_runner.py对每个已完成回合计算episode_length步数、episode_return累积回报、episode_duration_s耗时如果该回合此前被切成过 chunk还会把_ongoing_episodes_for_metrics中缓存的 chunk 数据累加进来保证统计的是完整回合通过_log_episode_metrics记录均值、最小 / 最大值等指标例如EPISODE_LEN_MEAN、EPISODE_RETURN_MEAN、EPISODE_RETURN_MIN/MAX等其平滑窗口按metrics_num_episodes_for_smoothing / num_env_runners折算以匹配算法进程中的并行合并逻辑single_agent_env_runner.py最终调用self.metrics.reduce()返回归约后的指标字典。同时_increase_sampled_metricssingle_agent_env_runner.py负责累计NUM_ENV_STEPS_SAMPLED、NUM_AGENT_STEPS_SAMPLED、NUM_MODULE_STEPS_SAMPLED、NUM_EPISODES等周期指标以及对应的_LIFETIME生命周期指标后者带吞吐统计with_throughputTrue。清理stop 与资源释放参考文档Cleanup一节只列出一个方法stop()env_runner.py其职责是释放该 EnvRunner 占用的所有资源——例如当内部使用 gym.Env 时应确保调用其close()方法。SingleAgentEnvRunner.stop()正是通过self.env.close()关闭向量化环境single_agent_env_runner.py。基类同时定义了__del__析构钩子在 Actor 被删除时清理资源。另外值得留意的是基类构造函数中num_env_steps_sampled_lifetime记录了该实例生命周期内累计采样的环境步数这一计数会随get_state()/set_state()在故障恢复、状态同步过程中被传递保证统计不断档single_agent_env_runner.py。容错StepFailedRecreateEnvError 与子环境重启StepFailedRecreateEnvError定义在 python/ray/rllib/env/env_errors.py是一个PublicAPI(stabilityalpha)的异常类。它表达的特殊语义是环境 step 失败且该失败是预期内的EnvRunner 应当重置 / 重建环境后继续而不是把错误当作异常向外传播。触发路径_try_env_stepEnvRunner 在步进环境时并不直接调用env.step()而是走_try_env_step()env_runner.py正常路径在ENV_STEP_TIMER计时下调用self.env.step(actions)成功后递增尝试步数计数器异常路径任何Exception都会被捕获向指标写入一次NUM_ENV_STEP_FAILURES_LIFETIME累计若config.restart_failed_sub_environmentsTrue除非异常本身就是StepFailedRecreateEnvError此时不打印错误日志避免噪音否则打印原始错误日志随后调用self.make_env()重建环境并返回ENV_STEP_FAILURE哨兵值。调用方_sample循环看到该哨兵值后会丢弃当前已采集的部分数据、重置所有环境并重试该步若config.restart_failed_sub_environmentsFalse记录完整错误信息并提示用户可通过fault_tolerance(restart_failed_sub_environmentsTrue)开启自动重建最终以RuntimeError形式向上抛出。_try_env_reset()采用同样的容错模式env_runner.py重置失败时若开启了restart_failed_sub_environments则重建环境后递归重试否则抛出原始异常。配置入口AlgorithmConfig.fault_tolerance()restart_failed_sub_environments的默认值为Falsealgorithm_config.py通过fault_tolerance()方法开启algorithm_config.pyfrom ray.rllib.algorithms.algorithm_config import AlgorithmConfig config ( AlgorithmConfig() .environment(CartPole-v1) .fault_tolerance( restart_failed_sub_environmentsTrue, # step 失败时静默重启子环境 ) )该方法同时提供restart_failed_env_runnersEnvRunner Actor 崩溃后以相同worker_index重建副本适用于可能被抢占的 SPOT 实例、ignore_env_runner_failures、max_num_env_runner_restarts、delay_between_env_runner_restarts_s、num_consecutive_env_runner_failures_tolerance、env_runner_health_probe_timeout_s、env_runner_restore_timeout_s等完整容错参数。需要特别区分的是restart_failed_sub_environments处理的是向量化环境内部某个子环境的失败属于 EnvRunner 内部的静默自愈不会影响 EnvRunner 本身而restart_failed_env_runners处理的是整个 EnvRunner Actor 的失败。关于StepFailedRecreateEnvError的使用场景源码注释给出了很好的指引当你的环境不稳定、会以某种固定模式崩溃时例如连接了一个自己无法完全控制的外部模拟器可以在step()方法中检测到这类崩溃后主动抛出该异常从而避免打印误导性的错误日志。同时注释也提醒谨慎使用如果环境持续崩溃该机制可能导致环境反复重置的无限循环env_errors.py。单智能体与多智能体 EnvRunner 的选择机制参考文档明确指出默认情况下 RLlib 提供两个内置的 EnvRunner 子类——面向单智能体的SingleAgentEnvRunner和面向多智能体的MultiAgentEnvRunner并根据你的配置自动决定使用哪一个。判断依据就是config.is_multi_agent属性。SingleAgentEnvRunnerSingleAgentEnvRunnerpython/ray/rllib/env/single_agent_env_runner.py用于单智能体场景继承EnvRunner并实现Checkpointable。其特点包括环境是 gymnasium 的gym.vector.VectorEnv经DictInfoToList包装单智能体RLModule负责前向推理构造时按条件worker_index is None、worker_index 0、create_env_on_local_worker、或num_env_runners 0决定本地 worker 是否创建环境构建 Env-to-Module 与 Module-to-Env 两套连接器connector管线负责观测预处理与动作后处理sample()返回List[SingleAgentEpisode]single_agent_episode.py通过get_state/set_state支持 checkpoint 与权重同步并兼容get_weights/set_weights的废弃告警。MultiAgentEnvRunnerMultiAgentEnvRunnerpython/ray/rllib/env/multi_agent_env_runner.py用于多智能体场景构造时若config.is_multi_agent为 False 会直接抛ValueError提示应通过config.multi_agent(policies..., policy_mapping_fn...)配置多智能体信息。其特点包括环境类型为VectorMultiAgentEnv模块为MultiRLModule按 ModuleID 管理多个子模块sample()返回List[MultiAgentEpisode]multi_agent_episode.py_sample循环中按config.count_steps_byenv_steps或agent_steps决定采样计步口径随机动作模式下只为本轮有观测、需要行动的智能体采样动作multi_agent_env_runner.py指标统计细化为按智能体EPISODE_AGENT_RETURN_MEAN、EPISODE_AGENT_STEPS与按模块EPISODE_MODULE_RETURN_MEAN两个维度。如何在你的配置中判断在代码中只需检查config.is_multi_agent即可知道当前是哪种配置多智能体环境的完整搭建方式可参考 multi_agent_env.rst 以及 RLlib 多智能体环境设置文档rllib-multi-agent-environments-doc 相关章节。两个子类的独立参考文档分别为 single_agent_env_runner.rst 与 multi_agent_env_runner.rst。总结EnvRunner 在新 API 栈中的位置把上述内容串联起来一条完整的采样链路是Algorithm或 Learner持有config.is_multi_agent决定实例化SingleAgentEnvRunner还是MultiAgentEnvRunner→ 通过ray.remote(EnvRunner)创建 N 个并行 Actornum_env_runners控制数量→ 每个 Actor 在__init__中完成make_env/make_module/ 连接器构建 → 训练循环反复调用sample()获取Episode数据 → 回合指标经get_metrics()归约后参与训练统计 → 训练结束后stop()释放资源期间若子环境崩溃StepFailedRecreateEnvError配合fault_tolerance(restart_failed_sub_environmentsTrue)实现无感自愈。理解 EnvRunner 的这套 API 与实现是掌握 RLlib 新 API 栈采样系统、诊断采样性能问题以及为自定义分布式 RL 场景扩展数据采集能力的基础。快速索引关注点参考文件EnvRunner 基类与容错实现env_runner.py单智能体实现single_agent_env_runner.py多智能体实现multi_agent_env_runner.py容错异常定义env_errors.pyfault_tolerance 配置algorithm_config.pyAPI 参考文档env_runner.rst、single_agent_env_runner.rst、multi_agent_env_runner.rst【免费下载链接】rayRay is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads.项目地址: https://gitcode.com/gh_mirrors/ra/ray创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考