verl 中的 Detached Worker 实战:基于 Ray 的跨进程 Worker 复用与远程训练调用

verl 中的 Detached Worker 实战:基于 Ray 的跨进程 Worker 复用与远程训练调用 verl 中的 Detached Worker 实战基于 Ray 的跨进程 Worker 复用与远程训练调用【免费下载链接】verlverl/HybridFlow: A Flexible and Efficient RL Post-Training Framework项目地址: https://gitcode.com/GitHub_Trending/ve/verl导读在 verlHybridFlow的分布式训练框架中single_controller模块负责把模型训练、Rollout、Reward 等角色封装为 Ray ActorWorker并由WorkerGroup统一编排。Detached Worker分离式 Worker是其中一种特殊用法Server 进程负责创建并常驻一组训练 WorkerClient 进程则通过 Worker 名称挂载到同一组已存活的 Worker 上直接发起 RPC 调用从而实现训练服务的进程间解耦。本文以 tests/single_controller/detached_worker 目录下的完整示例为骨架结合verl/single_controller的源码实现讲解如何在单节点上启动 Ray 集群、以 Server 身份创建 Detached Worker、以 Client 身份远程调用train_model完成一次 Megatron 训练迭代并深入剖析from_detached、Dispatch 模式与命名规则等底层机制帮助读者掌握将 verl 训练组件服务化的完整实战方案。一、背景verl 的 single_controller 抽象与 Detached Worker 的定位verl 的single_controller子模块是整套分布式训练编排的核心它将分布式训练中的每个进程抽象为Worker并将一组 Worker 聚合成WorkerGroup进行统一管理。相关基础类定义在 verl/single_controller/base/worker_group.py 中ResourcePool负责跨节点资源池管理ClassWithInitArgs负责延迟构造参数WorkerGroup则提供 Worker 管理、存活检查和远程方法绑定能力Ray 后端的具体实现位于 verl/single_controller/ray/base.py包括RayResourcePool、RayClassWithInitArgs和RayWorkerGroup。常规用法下训练进程自己创建RayWorkerGroup并随之销毁。而Detached Worker提供了一种先建后挂的模式Server 进程创建带detachedTrue的 Worker这些 Worker 以lifetimedetached方式注册到 Ray 集群即使创建它们的 Python 进程退出Actor 仍然存活Client 进程不重新创建 Worker而是通过RayWorkerGroup.from_detached(...)按 Worker 名称如trainerTrainer_0:0重新获取句柄形成对同一组 Worker 的远程句柄进而直接调用其注册的远程方法。从源码结构看worker_group.pyWorkerGroup.__init__通过self._is_init_with_detached_workers resource_pool is None区分两条初始化路径有资源池时新建 Worker没有资源池resource_poolNone时走_init_with_detached_workers附加既有 Worker。这正是 Detached Worker 模式的核心分水岭。二、示例总体架构Server/Client 分离的三步流程仓库中的示例位于tests/single_controller/detached_worker/目录共 4 个文件文件角色作用README.md使用说明三步运行流程server.pyServer启动 TrainerMegatron Llama 模型创建 Detached Workerclient.pyClient挂载已有 Worker构造数据发起远程训练run.sh一键脚本串联 ray start / server / client / ray stop按照 README.md 的说明标准运行方式分为三步目前仅支持单节点# 1. 启动本地 Ray 集群 ray start --head --port6379 # 2. 启动 Server创建 Detached Worker 并完成模型初始化 python3 server.py # 3. 在另一个终端运行 Client挂载 Worker 并发送训练数据 python3 client.py三个命令缺一不可没有 Ray 集群Actor 无处安放没有 Server 先启动Worker 尚未创建Client 将无法按名称找到目标Client 则是真正发起训练请求的一方。整个过程体现了服务端提供计算能力、客户端发起计算任务的 RPC 式协作模型。三、第一步启动本地 Ray 集群Detached Worker 的生存期绑定在 Ray 集群上因此必须先有集群。ray start --head --port6379会在当前节点启动一个 head 节点含 GCS、Driver 等使用默认端口 6379ray start --head --port6379需要注意的几点端口--port6379指定 GCSGlobal Control Store端口Client/Server 后续通过ray.init(addressauto)自动发现该集群单节点限制README 明确标注 Only on a single node这是当前示例的实现边界两个进程都通过ray.init(addressauto, namespaceverl)连入同一节点的同一个集群namespace 一致性Server 与 Client 都显式指定namespaceverl这是双方能在同一命名空间内互相发现 Actor 的前提。如果命名空间不一致ray.get_actor(name...)将找不到 Server 创建的 Trainer。运行结束后可用ray stop --force清理集群见 run.sh。四、第二步Server 端——创建 Detached Worker 组并初始化模型Server 的核心代码位于 server.py它定义了Trainer这一 Ray Actor 类并以 Detached 方式创建包含 2 个 Worker 的RayWorkerGroup。4.1 定义 Trainer ActorTrainer继承自 verl 的Worker基类verl/single_controller/base/worker.py并用ray.remote装饰使其成为 Ray Actorray.remote class Trainer(Worker): def __init__(self): super().__init__() if not torch.distributed.is_initialized(): rank int(os.environ[LOCAL_RANK]) torch.distributed.init_process_group(backendnccl) torch.cuda.set_device(rank) mpu.initialize_model_parallel( tensor_model_parallel_size2, pipeline_model_parallel_size1, virtual_pipeline_model_parallel_sizeNone, use_sharpFalse, context_parallel_size1, expert_model_parallel_size1, nccl_communicator_config_pathNone, ) tensor_parallel.model_parallel_cuda_manual_seed(10) is_collect ( mpu.get_tensor_model_parallel_rank() 0 and mpu.get_pipeline_model_parallel_rank() mpu.get_pipeline_model_parallel_world_size() - 1 and mpu.get_context_parallel_rank() 0 ) self._register_dispatch_collect_info( mesh_nametrain, dp_rankmpu.get_data_parallel_rank(), is_collectis_collect )关键点解读每个 Worker 在构造时初始化torch.distributed进程组并调用 Megatron 的initialize_model_parallel将 2 个 Worker 组织为tensor_model_parallel_size2的张量并行组_register_dispatch_collect_info(mesh_nametrain, dp_rank..., is_collect...)是 verl 的分布式编排钩子它将每个 Worker 在名为train的通信域中的 DP rank 与是否参与结果收集信息注册到 Worker 内部worker.py供后续 NDN-DimensionalDispatch 机制查询使用。4.2 用 register 声明远程方法Trainer 对外暴露两个远程方法分别使用不同的 Dispatch 模式decorator.pyregister(dispatch_modeDispatch.ONE_TO_ALL) def init_model(self): # 构造 LlamaConfigvocab_size256, hidden_size2048, ... # 通过 get_model(...) 构建 ParallelLlamaForCausalLMRmPadPP # 创建 Megatron Optimizerlr1e-6, clip_grad1.0 ... register(dispatch_modemake_nd_compute_dataproto_dispatch_fn(mesh_nametrain)) def train_model(self, data: DataProto) - DataProto: input_ids data.batch[input_ids] attention_mask data.batch[attention_mask] position_ids data.batch[position_ids] self.optimizer.zero_grad() self.model.zero_grad_buffer(...) output self.model(input_idsinput_ids, attention_maskattention_mask, position_idsposition_ids).logits output.mean().backward() update_successful, grad_norm, num_zeros_in_grad self.optimizer.step( self.megatron_config, self.megatron_config.timers ) return DataProto(batchTensorDict({loss: output.detach()}, batch_sizeoutput.shape[0]))Dispatch.ONE_TO_ALL把同一份参数广播到所有 Worker 上执行dispatch_one_to_all将每个参数复制world_size份见 decorator.py适合init_model这类每个进程各自初始化同一模型的操作make_nd_compute_dataproto_dispatch_fn(mesh_nametrain)返回一组dispatch_fn/collect_fn闭包decorator.py内部通过dispatch_lazy_compute_data_proto在运行时向 Worker 查询train网格的 DP rank 映射把DataProto按 DP 维度切分下发到对应 Worker训练结束后按is_collect掩码汇聚结果并concat回完整的DataProtodecorator.py。这解释了为什么train_model的输入输出都是DataProto它天然支持跨 TP/PP/CP 网格的数据分发与收集。4.3 以 Detached 方式创建 WorkerGroupServer 的__main__是 Detached Worker 模式的关键所在if __name__ __main__: ray.init(addressauto, namespaceverl) resource_pool RayResourcePool(process_on_nodes[2], detachedTrue) cls_with_init_args RayClassWithInitArgs(clsTrainer) worker_group RayWorkerGroup( resource_poolresource_pool, ray_cls_with_initcls_with_init_args, name_prefixtrainer, detachedTrue, ) worker_group.init_model() worker_names worker_group.worker_names print(worker_names)逐项说明RayResourcePool(process_on_nodes[2], detachedTrue)申请一个资源池在本节点上分配 2 个进程对应 2 个 GPUdetachedTrue使 Placement Group 也以 detached 生命周期创建ray/base.pyRayWorkerGroup(..., name_prefixtrainer, detachedTrue)name_prefix会作为 Actor 命名的前缀detachedTrue时在创建 Worker 时追加ray_cls_with_init.update_options({lifetime: detached})ray/base.py这是 Actor 脱离创建进程存活的直接原因worker_group.init_model()同步调用所有 Worker 的init_modelONE_TO_ALL 广播此时 Megatron 模型与优化器已就绪worker_group.worker_names打印出 Worker 名称供 Client 挂载使用。按 ray/base.py 的命名规则f{self.name_prefix}{cia_name}_{pg_idx}:{local_rank}此处输出应为类似trainerTrainer_0:0与trainerTrainer_0:1的两个名称——这正是 Client 端worker_names硬编码的来源。五、第三步Client 端——通过 from_detached 挂载并远程训练Client 的核心代码位于 client.py它不创建任何新 Worker而是借用Server 已创建好的那一组if __name__ __main__: ray.init(addressauto, namespaceverl) # get the worker group using names worker_names [trainerTrainer_0:0, trainerTrainer_0:1] cls_with_init_args RayClassWithInitArgs(clsTrainer) worker_group RayWorkerGroup.from_detached(worker_namesworker_names, ray_cls_with_initcls_with_init_args) batch_size 16 sequence_length 1024 # give Trainer some data to train input_ids torch.randint(low0, high256, size(batch_size, sequence_length), dtypetorch.int64, devicecuda) attention_mask torch.ones_like(input_ids) position_ids compute_position_id_with_mask(attention_mask) data DataProto( batchTensorDict( {input_ids: input_ids, attention_mask: attention_mask, position_ids: position_ids}, batch_sizebatch_size, ), meta_info{}, ) output worker_group.train_model(data) print(output)5.1 from_detached 的底层实现RayWorkerGroup.from_detached是类方法ray/base.py其本质是不携带资源池地构造一个 WorkerGroupclassmethod def from_detached(cls, name_prefixNone, worker_namesNone, worker_handlesNone, ray_cls_with_initNone, **kwargs): worker_group cls( resource_poolNone, # 关键resource_pool 置空 ray_cls_with_initray_cls_with_init, name_prefixname_prefix, worker_namesworker_names, worker_handlesworker_handles, **kwargs, ) return worker_group由于resource_poolNoneWorkerGroup.__init__判定_is_init_with_detached_workersTrue随即走_init_with_detached_workers路径ray/base.pydef _init_with_detached_workers(self, worker_names, worker_handles): # ray.get_actor holds a weak reference to the actor, which causes actors garbage collected unexpectedly # if we only hold spawn RayWorkerGroup. By passing actor handle explicitly, spawn RayWorkerGroup have # strong reference to these actors. workers worker_handles if worker_handles else [ray.get_actor(namename) for name in worker_names] self._workers workers self._world_size len(workers)即Client 通过ray.get_actor(nametrainerTrainer_0:0)等调用从 Ray 集群中按名称解析出 Server 创建的 Actor 句柄重建一个世界大小为 2 的 WorkerGroup。源码注释还提示了一个工程细节ray.get_actor持有的是弱引用若只持有从 spawn 出来的 WorkerGroup 而显式传入worker_handles可避免 Actor 被意外 GC。5.2 数据构造与远程训练调用Client 构造了一批随机数据batch_size16, sequence_length1024词表大小 256 与 Server 端LlamaConfig(vocab_size256)严格对应封装成DataProto后调用worker_group.train_model(data)。由于train_model注册的是 ND Compute DataProto 模式该调用会自动向各 Worker 查询train网格的 DP rank第一次调用时缓存到_dispatch_info将DataProto按 DP 维度切分下发在各 Worker 上执行一次前向 反向 Optimizer.stepMegatron 优化器lr1e-6, clip_grad1.0在is_collectTrue的 Worker 上收集结果并concat成完整DataProto返回输出为loss。最终 Client 打印出训练后的loss输出一次完整的远程训练迭代即告完成。需要注意的是由于是并行随机初始化本示例中模型并未经过预热训练loss 输出主要用于验证调用链路与数据流是否正确而非追求收敛效果。六、一键运行run.sh 串联全流程仓库同时提供了 run.sh将四步操作串成一条命令#!/bin/bash ray start --head --port6379 python3 server.py python3 client.py ray stop --force执行顺序说明ray start拉起集群python3 server.pyServer 前台运行创建 Detached Worker、初始化模型并打印 worker_names 后进程退出Detached Actor 不会随之消亡python3 client.pyClient 挂载 Worker、发起训练并打印输出ray stop --force收尾清理集群连带销毁所有 Detached Actor。该脚本把 README 中另一个终端手动执行 Client的步骤简化为顺序执行适合作为 CI 或本地快速验证的最小复现入口。七、原理深挖Detached Worker 的生命周期与命名机制7.1 双层 detached 语义Detached 在 verl 的 Ray 后端中体现为两层Placement Group 层RayResourcePool(detachedTrue)使 PG 的lifetimedetachedray/base.py资源组不会随创建进程退出而回收Actor 层RayWorkerGroup(detachedTrue)触发ray_cls_with_init.update_options({lifetime: detached})ray/base.pyActor 注册为 detached创建它的 Python 进程退出后仍由 Ray GCS 维护其存活。两层缺一不可仅 PG detached 而 Actor 非 detached进程退出后 Actor 仍可能被回收仅 Actor detached 而 PG 非 detached资源归属也可能在进程退出后失效。示例中 server.py 对RayResourcePool与RayWorkerGroup同时传入detachedTrue正是为了保证 Worker 组的完整持久化。7.2 Worker 名称的生成规则Client 端硬编码的worker_names并非随机而是严格遵循 ray/base.py 的命名模板cia_name type(ray_cls_with_init.cls).__name__ # 从 ActorClass(Trainer) 提取 Trainer name f{self.name_prefix}{cia_name}_{pg_idx}:{local_rank} # e.g. Worker_2:5代入本例name_prefixtrainer 类名Trainer 第 0 个 placement grouppg_idx0 局部 rank 0/1即得到trainerTrainer_0:0与trainerTrainer_0:1。理解这一规则对排查Client 找不到 Worker问题至关重要——名称对不上时ray.get_actor会直接抛错。7.3 Dispatch 模式与数据流Detached Worker 模式本身不引入新的 Dispatch 机制它复用的是single_controller标准的register装饰器体系decorator.py。预定义模式包括RANK_ZERO、ONE_TO_ALL、ALL_TO_ALL、DP_COMPUTE、DP_COMPUTE_PROTO等decorator.py而示例中用到的make_nd_compute_dataproto_dispatch_fn则是面向 TP/PP/CP 网格的通用数据分发方案。register会把dispatch_mode、execute_mode、blocking等属性以MAGIC_ATTR形式挂到方法上decorator.pyWorkerGroup初始化时通过func_generatorray/base.py将 Worker 的每个方法包装为Dispatch → 远程执行 → Collect三步的 Functorfrom_detached重建的 WorkerGroup 同样会执行_bind_worker_method因此 Client 拿到的train_model与 Server 自建时行为完全一致。八、扩展视角Detached Worker 在 verl 中的工程价值从源码结构看Detached Worker 机制是 verl 将训练组件服务化的基础设施其价值体现在进程角色解耦模型常驻Server与数据生产/调度Client分离可独立重启、独立扩缩容。示例中 Server 退出后 Worker 依然存活正是这种解耦的最小验证与spawn/fuse的互补RayWorkerGroup.spawn内部也调用from_detachedray/base.py来生成子角色 WorkerGroup 并重绑定方法前缀说明挂载机制是角色复用如 actor/critic/ref/rollout 拆分的公共底层能力面向真实训练的映射真实训练中RayWorkerGroup的name_prefix与worker_names由训练框架统一生成与管理可参考 trainer/ppo 与 workers 目录Detached 示例相当于把这一机制手工拆解为两个独立进程便于学习与调试。实践建议严格遵循三步执行顺序并保持 Server 与 Client 的namespace一致示例统一为verlClient 的worker_names应与 Server 打印的名称保持一致若修改了name_prefix或 Worker 数量process_on_nodes需同步更新本示例依赖 2 块 GPU 且仅支持单节点在真实集群中使用时需按 docs/start/install.rst 与 docs/start/multinode.rst 的指引准备环境并将 Detached 语义与资源池管理结合使用调试完成后及时执行ray stop --force避免 Detached Actor 与 Placement Group 长期占用 GPU 资源。九、小结本文以 tests/single_controller/detached_worker 为入口完整走通了 verl Detached Worker 的三种运行方式手动三步 / 一键脚本并溯源到 ray/base.py、decorator.py 与 worker.py 等核心源码梳理出双层 detached 生命周期、Actor 命名规则、from_detached挂载机制与 ND DataProto 分发链路。对希望将 verl 训练组件封装为常驻服务、或想深入理解single_controller编排原理的开发者而言这组示例是官方仓库中最直接的入门素材——先跑通示例再对照源码逐行阅读即可快速掌握 verl 分布式调度的核心心智模型。【免费下载链接】verlverl/HybridFlow: A Flexible and Efficient RL Post-Training Framework项目地址: https://gitcode.com/GitHub_Trending/ve/verl创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考