verl 中的 Detached Worker 实战:基于 Ray 的跨进程 Worker 复用与远程训练调用
2026/9/13 17:47:53 网站建设 项目流程

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

【免费下载链接】verlverl/HybridFlow: A Flexible and Efficient RL Post-Training Framework项目地址: https://gitcode.com/GitHub_Trending/ve/verl

导读

在 verl(HybridFlow)的分布式训练框架中,single_controller模块负责把模型训练、Rollout、Reward 等角色封装为 Ray Actor(Worker),并由WorkerGroup统一编排。Detached Worker(分离式 Worker)是其中一种特殊用法:Server 进程负责创建并常驻一组训练 Worker,Client 进程则通过 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,包括RayResourcePoolRayClassWithInitArgsRayWorkerGroup

常规用法下,训练进程自己创建RayWorkerGroup并随之销毁。而Detached Worker提供了一种"先建后挂"的模式:

  1. Server 进程:创建带detached=True的 Worker,这些 Worker 以lifetime="detached"方式注册到 Ray 集群,即使创建它们的 Python 进程退出,Actor 仍然存活;
  2. Client 进程:不重新创建 Worker,而是通过RayWorkerGroup.from_detached(...)按 Worker 名称(如trainerTrainer_0:0)重新获取句柄,形成对同一组 Worker 的"远程句柄",进而直接调用其注册的远程方法。

从源码结构看(worker_group.py),WorkerGroup.__init__通过self._is_init_with_detached_workers = resource_pool is None区分两条初始化路径:有资源池时新建 Worker,没有资源池(resource_pool=None)时走_init_with_detached_workers附加既有 Worker。这正是 Detached Worker 模式的核心分水岭。

二、示例总体架构:Server/Client 分离的三步流程

仓库中的示例位于tests/single_controller/detached_worker/目录,共 4 个文件:

文件角色作用
README.md使用说明三步运行流程
server.pyServer启动 Trainer(Megatron Llama 模型),创建 Detached Worker
client.pyClient挂载已有 Worker,构造数据发起远程训练
run.sh一键脚本串联 ray start / server / client / ray stop

按照 README.md 的说明,标准运行方式分为三步(目前仅支持单节点):

# 1. 启动本地 Ray 集群 ray start --head --port=6379 # 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 --port=6379会在当前节点启动一个 head 节点(含 GCS、Driver 等),使用默认端口 6379:

ray start --head --port=6379

需要注意的几点:

  • 端口--port=6379指定 GCS(Global Control Store)端口,Client/Server 后续通过ray.init(address="auto")自动发现该集群;
  • 单节点限制:README 明确标注 "Only on a single node",这是当前示例的实现边界(两个进程都通过ray.init(address="auto", namespace="verl")连入同一节点的同一个集群);
  • namespace 一致性:Server 与 Client 都显式指定namespace="verl",这是双方能在同一命名空间内互相发现 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 Actor

Trainer继承自 verl 的Worker基类(verl/single_controller/base/worker.py),并用@ray.remote装饰,使其成为 Ray Actor:

@ray.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(backend="nccl") torch.cuda.set_device(rank) mpu.initialize_model_parallel( tensor_model_parallel_size=2, pipeline_model_parallel_size=1, virtual_pipeline_model_parallel_size=None, use_sharp=False, context_parallel_size=1, expert_model_parallel_size=1, nccl_communicator_config_path=None, ) 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_name="train", dp_rank=mpu.get_data_parallel_rank(), is_collect=is_collect )

关键点解读:

  • 每个 Worker 在构造时初始化torch.distributed进程组,并调用 Megatron 的initialize_model_parallel,将 2 个 Worker 组织为tensor_model_parallel_size=2的张量并行组;
  • _register_dispatch_collect_info(mesh_name="train", dp_rank=..., is_collect=...)是 verl 的分布式编排钩子,它将每个 Worker 在名为"train"的通信域中的 DP rank 与"是否参与结果收集"信息注册到 Worker 内部(worker.py),供后续 ND(N-Dimensional)Dispatch 机制查询使用。

4.2 用 @register 声明远程方法

Trainer 对外暴露两个远程方法,分别使用不同的 Dispatch 模式(decorator.py):

@register(dispatch_mode=Dispatch.ONE_TO_ALL) def init_model(self): # 构造 LlamaConfig:vocab_size=256, hidden_size=2048, ... # 通过 get_model(...) 构建 ParallelLlamaForCausalLMRmPadPP # 创建 Megatron Optimizer(lr=1e-6, clip_grad=1.0) ... @register(dispatch_mode=make_nd_compute_dataproto_dispatch_fn(mesh_name="train")) 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_ids=input_ids, attention_mask=attention_mask, position_ids=position_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(batch=TensorDict({"loss": output.detach()}, batch_size=output.shape[0]))
  • Dispatch.ONE_TO_ALL:把同一份参数广播到所有 Worker 上执行(dispatch_one_to_all将每个参数复制world_size份,见 decorator.py),适合init_model这类"每个进程各自初始化同一模型"的操作;
  • make_nd_compute_dataproto_dispatch_fn(mesh_name="train"):返回一组dispatch_fn/collect_fn闭包(decorator.py),内部通过dispatch_lazy_compute_data_proto在运行时向 Worker 查询"train"网格的 DP rank 映射,把DataProto按 DP 维度切分下发到对应 Worker,训练结束后按is_collect掩码汇聚结果并concat回完整的DataProto(decorator.py)。这解释了为什么train_model的输入输出都是DataProto:它天然支持跨 TP/PP/CP 网格的数据分发与收集。

4.3 以 Detached 方式创建 WorkerGroup

Server 的__main__是 Detached Worker 模式的关键所在:

if __name__ == "__main__": ray.init(address="auto", namespace="verl") resource_pool = RayResourcePool(process_on_nodes=[2], detached=True) cls_with_init_args = RayClassWithInitArgs(cls=Trainer) worker_group = RayWorkerGroup( resource_pool=resource_pool, ray_cls_with_init=cls_with_init_args, name_prefix="trainer", detached=True, ) worker_group.init_model() worker_names = worker_group.worker_names print(worker_names)

逐项说明:

  • RayResourcePool(process_on_nodes=[2], detached=True):申请一个资源池,在本节点上分配 2 个进程(对应 2 个 GPU);detached=True使 Placement Group 也以 detached 生命周期创建(ray/base.py);
  • RayWorkerGroup(..., name_prefix="trainer", detached=True)name_prefix会作为 Actor 命名的前缀;detached=True时在创建 Worker 时追加ray_cls_with_init.update_options({"lifetime": "detached"})(ray/base.py),这是 Actor 脱离创建进程存活的直接原因;
  • worker_group.init_model():同步调用所有 Worker 的init_model(ONE_TO_ALL 广播),此时 Megatron 模型与优化器已就绪;
  • worker_group.worker_names打印出 Worker 名称,供 Client 挂载使用。按 ray/base.py 的命名规则f"{self.name_prefix}{cia_name}_{pg_idx}:{local_rank}",此处输出应为类似trainerTrainer_0:0trainerTrainer_0:1的两个名称——这正是 Client 端worker_names硬编码的来源。

五、第三步:Client 端——通过 from_detached 挂载并远程训练

Client 的核心代码位于 client.py,它不创建任何新 Worker,而是"借用"Server 已创建好的那一组:

if __name__ == "__main__": ray.init(address="auto", namespace="verl") # get the worker group using names worker_names = ["trainerTrainer_0:0", "trainerTrainer_0:1"] cls_with_init_args = RayClassWithInitArgs(cls=Trainer) worker_group = RayWorkerGroup.from_detached(worker_names=worker_names, ray_cls_with_init=cls_with_init_args) batch_size = 16 sequence_length = 1024 # give Trainer some data to train input_ids = torch.randint(low=0, high=256, size=(batch_size, sequence_length), dtype=torch.int64, device="cuda") attention_mask = torch.ones_like(input_ids) position_ids = compute_position_id_with_mask(attention_mask) data = DataProto( batch=TensorDict( {"input_ids": input_ids, "attention_mask": attention_mask, "position_ids": position_ids}, batch_size=batch_size, ), meta_info={}, ) output = worker_group.train_model(data) print(output)

5.1 from_detached 的底层实现

RayWorkerGroup.from_detached是类方法(ray/base.py),其本质是"不携带资源池地构造一个 WorkerGroup":

@classmethod def from_detached(cls, name_prefix=None, worker_names=None, worker_handles=None, ray_cls_with_init=None, **kwargs): worker_group = cls( resource_pool=None, # 关键:resource_pool 置空 ray_cls_with_init=ray_cls_with_init, name_prefix=name_prefix, worker_names=worker_names, worker_handles=worker_handles, **kwargs, ) return worker_group

由于resource_pool=NoneWorkerGroup.__init__判定_is_init_with_detached_workers=True,随即走_init_with_detached_workers路径(ray/base.py):

def _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(name=name) for name in worker_names] self._workers = workers self._world_size = len(workers)

即:Client 通过ray.get_actor(name="trainerTrainer_0:0")等调用,从 Ray 集群中按名称解析出 Server 创建的 Actor 句柄,重建一个世界大小为 2 的 WorkerGroup。源码注释还提示了一个工程细节:ray.get_actor持有的是弱引用,若只持有从 spawn 出来的 WorkerGroup 而显式传入worker_handles可避免 Actor 被意外 GC。

5.2 数据构造与远程训练调用

Client 构造了一批随机数据(batch_size=16, sequence_length=1024,词表大小 256 与 Server 端LlamaConfig(vocab_size=256)严格对应),封装成DataProto后调用worker_group.train_model(data)。由于train_model注册的是 ND Compute DataProto 模式,该调用会自动:

  1. 向各 Worker 查询"train"网格的 DP rank(第一次调用时缓存到_dispatch_info);
  2. DataProto按 DP 维度切分下发;
  3. 在各 Worker 上执行一次前向 + 反向 + Optimizer.step(Megatron 优化器,lr=1e-6, clip_grad=1.0);
  4. is_collect=True的 Worker 上收集结果并concat成完整DataProto返回(输出为loss)。

最终 Client 打印出训练后的loss输出,一次完整的"远程训练迭代"即告完成。需要注意的是:由于是并行随机初始化,本示例中模型并未经过预热训练,loss 输出主要用于验证调用链路与数据流是否正确,而非追求收敛效果。

六、一键运行:run.sh 串联全流程

仓库同时提供了 run.sh,将四步操作串成一条命令:

#!/bin/bash ray start --head --port=6379 python3 server.py python3 client.py ray stop --force

执行顺序说明:

  • ray start:拉起集群;
  • python3 server.py:Server 前台运行,创建 Detached Worker、初始化模型并打印 worker_names 后进程退出(Detached Actor 不会随之消亡);
  • python3 client.py:Client 挂载 Worker、发起训练并打印输出;
  • ray stop --force:收尾清理集群,连带销毁所有 Detached Actor。

该脚本把 README 中"另一个终端手动执行 Client"的步骤简化为顺序执行,适合作为 CI 或本地快速验证的最小复现入口。

七、原理深挖:Detached Worker 的生命周期与命名机制

7.1 双层 detached 语义

Detached 在 verl 的 Ray 后端中体现为两层:

  1. Placement Group 层RayResourcePool(detached=True)使 PG 的lifetime="detached"(ray/base.py),资源组不会随创建进程退出而回收;
  2. Actor 层RayWorkerGroup(detached=True)触发ray_cls_with_init.update_options({"lifetime": "detached"})(ray/base.py),Actor 注册为 detached,创建它的 Python 进程退出后仍由 Ray GCS 维护其存活。

两层缺一不可:仅 PG detached 而 Actor 非 detached,进程退出后 Actor 仍可能被回收;仅 Actor detached 而 PG 非 detached,资源归属也可能在进程退出后失效。示例中 server.py 对RayResourcePoolRayWorkerGroup同时传入detached=True,正是为了保证 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_prefix="trainer"+ 类名Trainer+ 第 0 个 placement group(pg_idx=0)+ 局部 rank 0/1,即得到trainerTrainer_0:0trainerTrainer_0:1。理解这一规则对排查"Client 找不到 Worker"问题至关重要——名称对不上时,ray.get_actor会直接抛错。

7.3 Dispatch 模式与数据流

Detached Worker 模式本身不引入新的 Dispatch 机制,它复用的是single_controller标准的@register装饰器体系(decorator.py)。预定义模式包括RANK_ZEROONE_TO_ALLALL_TO_ALLDP_COMPUTEDP_COMPUTE_PROTO等(decorator.py),而示例中用到的make_nd_compute_dataproto_dispatch_fn则是面向 TP/PP/CP 网格的通用数据分发方案。@register会把dispatch_modeexecute_modeblocking等属性以MAGIC_ATTR形式挂到方法上(decorator.py),WorkerGroup初始化时通过func_generator(ray/base.py)将 Worker 的每个方法包装为"Dispatch → 远程执行 → Collect"三步的 Functor,from_detached重建的 WorkerGroup 同样会执行_bind_worker_method,因此 Client 拿到的train_model与 Server 自建时行为完全一致。

八、扩展视角:Detached Worker 在 verl 中的工程价值

从源码结构看,Detached Worker 机制是 verl 将"训练组件服务化"的基础设施,其价值体现在:

  • 进程角色解耦:模型常驻(Server)与数据生产/调度(Client)分离,可独立重启、独立扩缩容。示例中 Server 退出后 Worker 依然存活,正是这种解耦的最小验证;
  • spawn/fuse的互补RayWorkerGroup.spawn内部也调用from_detached(ray/base.py)来生成子角色 WorkerGroup 并重绑定方法前缀,说明挂载机制是角色复用(如 actor/critic/ref/rollout 拆分)的公共底层能力;
  • 面向真实训练的映射:真实训练中RayWorkerGroupname_prefixworker_names由训练框架统一生成与管理(可参考 trainer/ppo 与 workers 目录),Detached 示例相当于把这一机制手工拆解为两个独立进程,便于学习与调试。

实践建议:

  • 严格遵循三步执行顺序,并保持 Server 与 Client 的namespace一致(示例统一为"verl");
  • Client 的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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询