跳转至

remote.py — RemoteRewardManager 将奖励计算分发到多个 Ray 远程 worker 上并行...

文件路径: verl/experimental/reward_loop/reward_manager/remote.py

文件概述

RemoteRewardManager 将奖励计算分发到多个 Ray 远程 worker 上并行执行。适用于 CPU 密集型的奖励计算场景(如代码执行评测、复杂规则验证),通过多进程并行加速。

关键代码

远程计算 Worker

@ray.remote
class RewardComputeWorker:
    """Ray 远程 worker,执行具体的奖励计算"""

    def __init__(self, compute_score_fn):
        self.compute_score_fn = compute_score_fn

    def compute(self, prompt, response, ...):
        """在远程进程中执行评分函数"""
        return self.compute_score_fn(prompt, response, ...)

@ray.remote 装饰器让这个类的实例运行在独立的进程中,Ray 负责进程间通信和任务调度。

RemoteRewardManager 类

@register("remote")
class RemoteRewardManager(RewardManagerBase):
    def __init__(self, config, tokenizer, num_workers=8, ...):
        super().__init__(config, tokenizer, ...)
        # 创建多个远程 worker
        self.workers = [
            RewardComputeWorker.remote(self.compute_score_fn)
            for _ in range(num_workers)
        ]

并行计算流程

def run_single(self, data, compute_score_fn, ...):
    # 将数据分发给各个 worker
    futures = []
    for i, (prompt, response) in enumerate(zip(prompts, responses)):
        worker = self.workers[i % len(self.workers)]  # 轮询分配
        future = worker.compute.remote(prompt, response, ...)
        futures.append(future)

    # 等待所有 worker 完成
    scores = ray.get(futures)

    return torch.tensor(scores, dtype=torch.float32)

工作原理: 1. 创建 N 个远程 worker(独立进程) 2. 将样本轮询分配给各 worker 3. 所有 worker 并行计算 4. 收集结果并返回

核心类/函数列表

名称 类型 说明
RewardComputeWorker Ray Actor 远程计算 worker
RemoteRewardManager 类 基于 Ray 的并行奖励管理器
num_workers 参数 远程 worker 数量

与其他模块的关系

  • 继承自 RewardManagerBase
  • 通过 @register("remote") 注册到全局注册表
  • 使用 Ray 框架进行分布式计算
  • 适用于 sandbox 代码执行、数学验证等 CPU 密集场景

小结

RemoteRewardManager 通过 Ray 的分布式计算能力,将奖励计算并行化。当评分函数是 CPU 密集型(不需要 GPU)但耗时较长时,这种并行方案比单线程快很多倍。