十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Ray RLlib Learner Connector 管道实战:从 Episode 到 Train Batch 的编译管线与自定义改造

Ray RLlib Learner Connector 管道实战:从 Episode 到 Train Batch 的编译管线与自定义改造 人工智能分布式训练强化学习任务调度模型推理服务【免费下载链接】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点击查看免费下载导读本文基于 Ray RLlib 的 learner-connector 文档 展开深入讲解 Learner ConnectorV2 管道Learner connector pipeline的工作机制每个 Learner actor 如何通过一条管道把一批 Episode 编译成可直接送入 RLModule.forward_train() 的 train batch。你将掌握默认管道的组成与顺序、通过 AlgorithmConfig 挂载自定义 ConnectorV2 片段的两种方式并学会实现损失计算前的内在奖励塑造reward shaping与最近 N 帧观测堆叠frame stacking两大实战场景同时结合仓库源码rllib/connectors/与官方示例count_based_curiosity.py、atari_ppo.py印证其底层原理。如上图所示Learner connector pipeline 位于输入训练数据一批 episodes与 Learner actor 的 RLModule 之间。它把输入数据转换成 RLModule.forward_train() 可直接读取的 tensor batchLearner actor 将管道输出直接送入该方法的forward_train()调用。一、Learner connector pipeline 是什么1.1 三类 ConnectorV2 管道中的训练端RLlib 的 ConnectorV2 抽象定义了connector piece连接片段的 API多个片段串成connector pipeline连接管道。从 ConnectorV2 的类文档可以看出任何 ConnectorV2 片段都可能属于以下三类管道之一EnvToModulePipeline位于 EnvRunner 侧把环境输出数据转换成 RLModule 可读的数据用于下一次forward_exploration()/forward_inference()前向计算典型职责包括观测后处理、观测过滤器、RNN 时间序列与零填充准备ModuleToEnvPipeline同样位于 EnvRunner 侧把 RLModule 输出如动作分布参数转换成可下发给env.step()的最终动作LearnerConnectorPipeline位于 Learner worker 侧把EnvRunner.sample()或回放缓冲区replay buffer产出的原始训练数据batch 或 episode 列表转换成 RLModuleforward_train()可读的训练数据用于损失计算。本文主角正是第三类 LearnerConnectorPipeline其实现类是 LearnerConnectorPipeline被标注为PublicAPI(stabilityalpha)。它继承自 ConnectorPipelineV2在__call__()中额外通过 MetricsLogger 记录进入/流出管道的 episode 长度总和LEARNER_CONNECTOR_SUM_EPISODES_LENGTH_IN/OUT用于观测管道对数据规模的影响同时它允许用户传入空 batchbatch{}由后续片段从 episodes 中自行填充。1.2 管道输入输出输入一批 Episode 对象单智能体SingleAgentEpisode或多智能体MultiAgentEpisode它们来自环境采样或回放缓冲区。输出一个RLModule可读的、按列column组织的 tensor batch如obs、actions、rewards、terminateds、truncateds、seq_lens、loss_mask、state_in等直接送入 RLModule.forward_train()。二、默认 Learner pipeline 行为默认情况下RLlib 会为每个 Learner connector pipeline 自动填充以下内置片段顺序即执行顺序。这一默认装配逻辑可以在 AlgorithmConfig.build_learner_connector() 中看到完整实现。顺序ConnectorV2 片段类路径职责1AddObservationsFromEpisodesToBatchrllib/connectors/common/add_observations_from_episodes_to_batch.py把传入 episodes 的所有观测放入 batch 的obs列2AddColumnsFromEpisodesToTrainBatchrllib/connectors/learner/add_columns_from_episodes_to_train_batch.py把其余列rewards、actions、termination 标志等放入 batch3AddTimeDimToBatchAndZeroPad仅状态模型相关rllib/connectors/common/add_time_dim_to_batch_and_zero_pad.py若 RLModule 是有状态stateful的在 axis1 增加大小为max_seq_len的时间维并在 episode 结束于不可被max_seq_len整除的时间步处右端零填充4AddStatesFromEpisodesToBatch仅状态模型相关rllib/connectors/common/add_states_from_episodes_to_batch.py若 RLModule 是有状态的把模块最近一次的状态输出作为新的状态输入放入 batch 的state_in列该列不含时间维5AgentToModuleMapping仅多智能体rllib/connectors/common/agent_to_module_mapping.py依据多智能体 episode 中已确定的 agent-to-module 映射把每个 agent 的数据映射到对应 module 的数据6BatchIndividualItemsrllib/connectors/common/batch_individual_items.py把 batch 中目前仍是逐条 item 列表的数据转换成批量结构即 axis0 为 batch 轴的 NumPy 数组7NumpyToTensorrllib/connectors/common/numpy_to_tensor.py把 batch 中的所有 NumPy 数组转换为框架张量并按需搬到 GPU2.1 默认片段细节源码级AddObservationsFromEpisodesToBatch作为 Learner connectoras_learner_connectorTrue时它会把每个 episode 的全部观测除最后一个终止观测外见源码 add_observations_from_episodes_to_batch.py 注释终止 episode 的最后一个观测对训练没有价值通过add_n_batch_items()加入 batch 的obs列。例如两个长度分别为 10 和 20 的 episodes最终 train batch 大小为 30。注意该片段不会改动 episodes 中的任何数据。AddColumnsFromEpisodesToTrainBatch补齐 actions、rewards、terminateds、truncateds 等列同时把 episode 中extra_model_outputs如 action logp、action probs放入对应列state_out除外它由AddStatesFromEpisodesToBatch处理。关键点若某列已在 batch 中用户自定义片段已填充则跳过、直接透传见源码 add_columns_from_episodes_to_train_batch.py 中if Columns.XXX not in batch的守卫逻辑。AddTimeDimToBatchAndZeroPad仅当rl_module.is_stateful()为 True 时生效。作为 Learner connector 时它会按model_config[max_seq_len]对每条数据执行 split_and_zero_pad()把长 episode 切成多个(B, T, ...)块并额外生成seq_lens与loss_mask列见源码 add_time_dim_to_batch_and_zero_pad.py供后续 RNN/Transformer 前向与损失屏蔽使用。若使用有状态模块却未提供max_seq_len该片段会抛出 ValueError提示通过config.rl_module(model_config{max_seq_len: ...})设置。AddStatesFromEpisodesToBatch对有状态模块把 episode 记录的最新state_out作为state_in放入 batch无时间维。从源码 add_states_from_episodes_to_batch.py 可看到其状态来源逻辑episode 从零开始t_started 0或没有state_out时使用rl_module.get_initial_state()否则使用上一 chunk 末尾的state_out作为回看状态。每max_seq_len步取一个状态最后的state_out被忽略与最后观测同理。2.2 关闭默认片段你可以通过如下配置一次性禁用上述所有默认片段config.learners(add_default_connectors_to_learner_pipelineFalse)对应配置项在 algorithm_config.py 中默认值为True源码 build_learner_connector() 中只有其为 True 时才追加默认片段。关闭后管道中仅保留你自定义的片段所有默认行为都由你自行实现。2.3 顺序的重要性这些变换的执行顺序对管道语义至关重要必须先有原始列数据才能做时间维切分与零填充必须先有切分后的数据才能取每段对应的state_inBatchIndividualItems批量化和NumpyToTensor张量化必须放在最后。这一点也体现在默认装配顺序与各片段文档中的 pipeline 注释如 numpy_to_tensor.py 中完整列出的默认管道结构中。三、编写自定义 Learner connector3.1 自定义函数的签名与挂载要在 Learner 连接管道中插入自定义逻辑需在 AlgorithmConfig 中指定一个函数该函数接收观测空间obs space和动作空间act space作为输入参数返回单个 ConnectorV2 片段或片段列表。config.learners( learner_connectorlambda obs_space, act_space: MyLearnerConnector(..), )需要添加多个自定义片段时以列表形式返回# 返回一个列表让 RLlib 把它们全部加入你的 Learner pipeline。 config.learners( learner_connectorlambda obs_space, act_space: [ MyLearnerConnector(..), MyOtherLearnerConnector(..), AndOneMoreConnector(..), ], )RLlib 会把函数返回的片段**按返回顺序前插prepend**到默认 Learner pipeline 之前如果配置了add_default_connectors_to_learner_pipelineFalse则仅使用你提供的片段、不加任何默认行为。从源码 build_learner_connector() 可以看到val_为单个ConnectorV2时包装成[val_]为 list/tuple 时直接展开否则抛出ValueError提示必须返回 ConnectorV2 对象或列表。兼容性提示config.training(learner_connector..)仍是支持的别名写法但在仓库源码 algorithm_config.py 中已被标记为 deprecated 并提示改用config.learners(learner_connector..)。自定义片段插入管道头部是刻意的设计如果你的自定义片段以任何方式修改了输入 episodes例如修改奖励管道末尾的默认片段会自动把这些修改后的数据加入 train batch。3.2 ConnectorV2 片段基类速览自定义片段需继承 ConnectorV2通常只需覆写两个方法__call__(self, *, rl_module, batch, episodes, exploreNone, shared_dataNone, **kwargs)核心变换逻辑。输入是前一个片段或空 dict的输出 batch、当前处理的所有 episodes 等返回新的 batch。签名使用 keyword-only 参数。recompute_output_observation_space(input_observation_space, input_action_space)可选当自定义片段改变了观测形状/类型时返回调整后的观测空间使 RLlib 能正确推导后续 RLModule 的输入空间。基类还提供了几个非常实用的工具方法可在 connector_v2.py 中查看完整实现single_agent_episode_iterator(episodes, agents_that_stepped_only...)把单/多智能体 episodes 统一展开成单智能体 episode 迭代器add_batch_item(batch, column, item_to_add, single_agent_episode)往 batch 的某列添加一个item自动按 episode 组织数据add_n_batch_items(batch, column, items_to_add, num_items, single_agent_episode)往 batch 的某列批量添加多个item。四、实战一损失计算前的奖励塑造内在奖励4.1 应用场景与动机写自定义 Learner ConnectorV2 片段的一个典型场景是在损失计算前进行奖励塑造reward shaping。Learner connector 的__call__()拥有对完整 episode 数据的全部访问权限观测、动作、多智能体场景下其他 agent 的数据、以及全部奖励。文档给出了一个简单而有效的计数型内在奖励count-based intrinsic reward信号自定义片段把内在奖励定义为agent 已经看到某个特定观测的次数的倒数。agent 访问某个状态越频繁该状态的内在奖励越低从而激励 agent 探索新状态、表现出更好的探索行为。4.2 完整代码实现仓库提供了完整可运行示例脚本rllib/examples/curiosity/count_based_curiosity.py。该示例在稀疏奖励环境 FrozenLake 8x8、时间步上限 14 上对比了使用与不使用 curiosity 的 PPO 策略只有使用计数型 curiosity 的策略能够真正学会任务脚本头注释给出了预期实验结果。运行方式python rllib/examples/curiosity/count_based_curiosity.py # 开启 curiosity python rllib/examples/curiosity/count_based_curiosity.py --no-curiosity # 关闭用于对照 # 调试模式便于设置断点 python rllib/examples/curiosity/count_based_curiosity.py --no-tune --num-env-runners0示例还支持--intrinsic-reward-coeff参数默认 1.0用于控制内在奖励叠加到外在奖励前的加权系数。首先通过继承 ConnectorV2 并覆写__call__编写自定义片段下述代码与文档一致也是 CountBasedCuriosity 的简化版from collections import Counter from ray.rllib.connectors.connector_v2 import ConnectorV2 class CountBasedIntrinsicRewards(ConnectorV2): def __init__(self, **kwargs): super().__init__(**kwargs) # 观测计数器用于计算状态访问频率。 self._counts Counter()在__call__中遍历所有单智能体 episode把其中存储的奖励改写为r(t) re(t) 1 / N(ot)其中re是环境给出的外在奖励N(ot)是 agent 到达观测o(t)的次数。def __call__( self, *, rl_module, batch, episodes, exploreNone, shared_dataNone, **kwargs, ): for sa_episode in self.single_agent_episode_iterator( episodesepisodes, agents_that_stepped_onlyFalse ): # 遍历除最后一个之外的所有观测。 observations sa_episode.get_observations(slice(None, -1)) # 获取所有对应的外在奖励。 rewards sa_episode.get_rewards() for i, (obs, rew) in enumerate(zip(observations, rewards)): # 计数器 1。 obs tuple(obs) self._counts[obs] 1 # 计算计数型内在奖励并叠加到外在奖励上。 rew 1 / self._counts[obs] # 把新奖励写回 episode正确的 timestep/索引处。 sa_episode.set_rewards(new_datarew, at_indicesi) return batch4.3 挂载与关键设计要点通过算法配置把自定义片段接入管道config.learners(learner_connectorlambda env: CountBasedIntrinsicRewards())挂载后损失函数即可在 incoming batch 的rewards列中收到被改写后的奖励信号。这里有一个关键设计点文档以 note 形式强调自定义逻辑把新奖励写回给定的 episodes而不是直接写 train batch。这样确保只有被修改的数据对后续片段可见batch 本身最初保持不变随后默认 Learner 片段之一 AddColumnsFromEpisodesToTrainBatch 会从 episodes 中提取奖励数据填充 batch因此你对 episode 对象做的任何修改都会被自动加入 train batch。这也再次印证了自定义片段前插在默认片段之前这一设计的意义。五、实战二堆叠最近 N 帧观测Frame Stacking5.1 为什么用两个管道配合实现帧堆叠Learner connector API 的另一个典型应用场景是结合自定义 env-to-module connector 实现高效的观测帧堆叠observation frame stacking。其优势在于无需对堆叠后相互重叠的观测数据去重也无需把这些额外观测存进 episodes更无需通过网络在 actor 之间传输这些数据。由于你没有覆写收集到的 episodes 中原始未堆叠的观测因此同一套 batch 构建逻辑必须执行两次一次在 EnvRunner actor 上用于动作计算一次在 Learner actor 上用于损失计算。其根本原因是管道产出的 batch 是临时的ephemeralRLlib 在 RLModule 前向结束后立即丢弃帧堆叠直接作用于正在构建的 batch 上因为你不想让去重后的堆叠观测塞满 episodes。5.2 一个类同时覆盖两个管道文档给出了一个StackFourObservations示例单个 ConnectorV2 类即可同时承担 env-to-module 与 Learner 两端的自定义片段职责。它以观测为 1D 张量的环境为例实现堆叠最近四帧观测import gymnasium as gym import numpy as np from ray.rllib.connectors.connector_v2 import ConnectorV2 from ray.rllib.core.columns import Columns class StackFourObservations(ConnectorV2): 一个把最近四个观测堆叠为一个的 connector 片段。 既可以作为 Learner connector也可以作为 env-to-module connector。 def recompute_output_observation_space( self, input_observation_space, input_action_space, ): # 假定输入观测空间是形状为 (x,) 的 Box。 assert ( isinstance(input_observation_space, gym.spaces.Box) and len(input_observation_space.shape) 1 ) # 该 connector 在 axis0 上拼接最近四个观测因此输出空间形状为 (4*x,)。 return gym.spaces.Box( lowinput_observation_space.low, highinput_observation_space.high, shape(input_observation_space.shape[0] * 4,), dtypeinput_observation_space.dtype, ) def __init__( self, input_observation_space, input_action_space, *, as_learner_connector, **kwargs, ): super().__init__(input_observation_space, input_action_space, **kwargs) self._as_learner_connector as_learner_connector def __call__(self, *, rl_module, batch, episodes, **kwargs): # 遍历所有单智能体episode。 for sa_episode in self.single_agent_episode_iterator(episodes): # 从 episode 中取最近四个观测。 last_4_obs sa_episode.get_observations( indices[-4, -3, -2, -1], fill0.0, # 到达 episode 开头时左端零填充。 ) # 拼接所有被堆叠的观测。 new_obs np.concatenate(last_4_obs, axis0) # 使用 ConnectorV2.add_batch_item() 工具方法把堆叠观测加入 batch。 # 注意这里不修改 episode。这意味着如果 self 是 env-to-module # connector而非 Learner connector采集到的 episode 中仍然只有 # 单个、未堆叠的观测Learner 管道必须为 forward_train() 前向再次堆叠。 self.add_batch_item( batchbatch, columnColumns.OBS, item_to_addnew_obs, single_agent_episodesa_episode, ) # 返回含堆叠观测的batch。 return batch然后在 AlgorithmConfig 中添加如下两行配置分别作用于两端管道from ray.rllib.algorithms.ppo import PPOConfig config PPOConfig() # 在 EnvRunner 侧启用帧堆叠。 config.env_runners( env_to_module_connectorlambda env, spaces, device: StackFourObservations(), ) # 在 Learner 侧再次启用帧堆叠。 config.training( learner_connectorlambda obs_space, act_space: StackFourObservations( as_learner_connectorTrue ), )5.3 观测空间自动推导你的 RLModule 会在其 setup() 方法中自动收到正确调整后的观测空间。EnvRunner 及其 env-to-module connector pipeline 会通过各片段的 recompute_output_observation_space() 方法为你计算这一信息该方法把片段串联的观测空间变换逐级传导。务必确保你的 RLModule 支持堆叠后的观测而不是单个观测。另外你不必像上述实现那样把观测在原有维度上拼接——只要 RLModule 知道如何处理改变后的观测形状你也可以把帧堆叠到一个全新的观测维度上。5.4 现成的 FrameStacking 片段推荐上文代码仅用于演示。RLlib 已内置一个开箱即用的 FrameStacking ConnectorV2 片段可同时用于 env-to-module 与 Learner connector 管道且支持多智能体场景。添加以下配置即可开启观测帧堆叠from ray.rllib.connectors.common.frame_stacking import FrameStacking N 4 # 要堆叠的帧数 # EnvRunner 侧帧堆叠。 config.env_runners( env_to_module_connectorlambda env, spaces, device: FrameStacking(num_framesN), ) # Learner 侧再次帧堆叠。 config.training( learner_connectorlambda obs_space, act_space: FrameStacking(num_framesN, as_learner_connectorTrue), )从 frame_stacking.py 源码可见其参数num_frames堆叠帧数默认 1、multi_agent是否处理多智能体观测空间、as_learner_connector是否作为 Learner 片段。Learner 侧实现通过np.lib.stride_tricks.as_strided构造滑动窗口视图并转置为(B, ..., num_frames)通道布局见源码 frame_stacking.py要求观测空间为Box且最后一维为 1见_convert_individual_space中的断言即通道位于最后一维的布局。仓库还提供了完整的端到端 Atari 示例 rllib/examples/algorithms/ppo/atari_ppo.py它在 Pong 上以FrameStackingEnvToModule(num_frames4)和FrameStackingLearner(num_frames4)分别挂载两端管道见脚本第 36-41 行并在_env_creator中通过wrap_atari_for_new_api_stack(..., framestackNone)关闭 RLlib 硬编码的帧堆叠改由 ConnectorV2 API 完成第 46-51 行运行命令参考脚本顶部注释python rllib/examples/algorithms/ppo/atari_ppo.py六、调试与指标观测Learner connector 管道内置了指标记录能力LearnerConnectorPipeline.call() 会在管道执行前后通过 MetricsLogger 记录进入与流出管道的所有 episode 长度总和LEARNER_CONNECTOR_SUM_EPISODES_LENGTH_IN与LEARNER_CONNECTOR_SUM_EPISODES_LENGTH_OUT。两者之差可用于诊断自定义片段是否意外增删了 episode 数据例如未正确执行sa_episode.set_rewards()而只是修改了 batch。总结Learner connector pipeline 是 Ray RLlib 新 API 栈中连接episode 原始数据与RLModule 损失计算的标准化编译管线。理解它的默认装配观测/列提取 → 时间维与零填充 → 状态输入 → 多智能体映射 → 批量聚合 → 张量化与自定义片段前插的扩展机制可以让你在不侵入算法核心逻辑的前提下通过 config.learners(learner_connector...) 优雅地实现奖励塑造、观测帧堆叠、多智能体数据重组等训练端定制。建议进一步阅读ConnectorV2 概览三类管道与片段 API 的完整说明Env-to-module connector 文档动作计算侧对应管道RLlib 新 API 栈迁移指南config.learners()/config.env_runners()新配置入口源码rllib/connectors/ 下所有默认片段实现与其内嵌测试代码是学习 ConnectorV2 行为的最佳素材。赞分享人工智能分布式训练强化学习任务调度模型推理服务【免费下载链接】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点击查看免费下载相关推荐Driver Store Explorer彻底告别Windows驱动臃肿轻松释放数GB系统盘空间Driver Store Explorer彻底告别Windows驱动臃肿轻松释放数GB系统盘空间 你是不是经常发现Windows系统盘空间在不知不觉中减少人工智能分布式训练强化学习任务调度模型推理服务Ray Train 分布式 PyTorch 回归训练实战从 CSV 数据管道到 TorchTrainer 全解析Ray Train 分布式 PyTorch 回归训练实战从 CSV 数据管道到 TorchTrainer 全解析 导读 本文围绕 doc/source/tra人工智能分布式训练强化学习任务调度模型推理服务Ray项目中的RLlib单智能体Episode详解Ray项目中的RLlib单智能体Episode详解 引言为什么需要深入理解Episode 在强化学习Reinforcement Learning实践中人工智能分布式训练强化学习任务调度模型推理服务上一篇Aurora博客系统从零开始搭建现代化个人博客的完整指南下一篇Mugen 项目推荐创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表