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)但耗时较长时,这种并行方案比单线程快很多倍。