Python Ray脚本设置num_cpus=4却仅在单worker运行的问题排查
Ray设置num_cpus=4但任务始终单Worker运行问题排查
问题复现代码
import ray import time import h5py @ray.remote class Analysis: def __init__(self): self._file = h5py.File('./Data/Trajectories/MDANSE/apoferritin.h5') def __getstate__(self): print('I dump') d = self.__dict__.copy() del d['_file'] return d def __setstate__(self,state): self.__dict__ = state self._file = h5py.File('./Data/Trajectories/MDANSE/apoferritin.h5') def run_step(self,index): time.sleep(5) print('I run a step',index) def combine(self,index): print('I combine',index) ray.init(num_cpus=4) a = Analysis.remote() obj_id = ray.put(a) for i in range(100): output = ray.get(a.run_step.remote(i))
上述代码初始化Ray时指定num_cpus=4,预期启动4个Worker并行执行任务,实际运行时所有任务始终调度到单个Worker上串行执行。
问题根因
代码存在3个直接导致无法并行的问题:
- 错误使用单Actor实例期望并行:
@ray.remote修饰的类是Ray Actor,遵循Actor模型设计,单个Actor实例只会固定绑定到一个Worker进程运行,所有提交给该实例的方法调用默认串行排队执行。代码中仅创建了1个Analysis实例a,所有run_step任务都提交给这一个实例,自然只能在单个Worker上运行,和CPU配额多少无关。 - 循环内同步阻塞提交任务:每次提交
a.run_step.remote(i)后立刻调用ray.get()阻塞等待当前任务执行完成,才会进入下一轮循环提交下一个任务。哪怕不用Actor、用普通无状态Remote函数,这种写法也会完全退化成串行执行,根本没有给Ray批量调度任务的机会。 - 冗余的
obj_id = ray.put(a)无实际作用:Actor实例化后返回的a本身就是指向远程对象的句柄,不需要额外调用ray.put()重复存入对象存储,这行代码不影响并行逻辑但属于无效写法。
修正方案
如果要实现4Worker并行执行任务,可按以下逻辑调整:
- 无状态计算场景优先用普通
@ray.remote函数,不要套Actor类;如果需要维护实例状态(比如持有h5文件句柄),则创建和CPU数量匹配的多个Actor实例,把任务分发给不同实例。 - 先批量提交所有任务拿到全部ObjectRef,最后统一调用
ray.get()等待所有结果返回,不要在提交循环里阻塞。
参考修正示例(无状态场景,用普通Remote函数实现并行):
import ray import time import h5py # 无状态逻辑直接用remote函数,自动调度到空闲worker @ray.remote def run_step(index): # 注意:如果每个任务都要读h5文件,建议每个worker初始化一次文件句柄,不要每次任务都重复打开 time.sleep(5) print('I run a step', index) return index ray.init(num_cpus=4) # 先批量提交所有任务,不要逐次get task_refs = [run_step.remote(i) for i in range(100)] # 最后统一等待所有任务完成 outputs = ray.get(task_refs)
内容的提问来源于stack exchange,提问作者Eurydice
相关产品推荐
相关产品推荐

