嵌套循环运行Ray任务触发内存不足问题的解决办法
解决Ray嵌套循环内存不足的方案
1. 控制并发任务总量,分批提交执行
原代码会一次性生成约1250万任务(5000行dataA的两两组合),远超系统内存承载上限。需要分批提交任务,限制同时运行的任务数量,避免任务队列和结果占用过多内存。
修改后的代码示例:
import ray import os from itertools import islice num_cpu = os.cpu_count() ray.init(num_cpus=num_cpu) # 提前将大对象存入Ray对象存储,避免重复传递 dataB_ref = ray.put(dataB) # 5000 x 256 dataC_ref = ray.put(dataC) # 5000 x 768 @ray.remote def calculate(dataB, dataC, i, j): # 仅处理i、j对应的子集,避免全量加载大矩阵 calcResult = someCalculationWithDataBAndC(dataB, dataC, i, j) return calcResult # 生成所有需要处理的(i,j)对 task_pairs = [] for i in range(len(dataA)): for j in range(i + 1, len(dataA)): task_pairs.append((i, j)) # 按CPU核心数的2倍设置批次大小(可根据内存情况调整) batch_size = num_cpu * 2 all_results = [] # 分批提交并处理任务 for batch in iter(lambda: list(islice(task_pairs, batch_size)), []): current_tasks = [calculate.remote(dataB_ref, dataC_ref, i, j) for i, j in batch] # 等待当前批次完成,获取结果后及时释放任务引用 batch_results = ray.get(current_tasks) all_results.extend(batch_results) ray.shutdown()
2. 优化任务内的内存占用
- 确保
someCalculationWithDataBAndC只读取dataB、dataC中与i、j相关的子集(比如numpy数组切片dataB[i]、dataC[j]),不要加载整个大矩阵到Worker进程内存。 - 如果任务不需要完整的dataB/dataC,可在Driver端提前拆分数据为对应i、j的小片段,再传递给任务(Ray对象存储会自动共享内存,切片操作不会额外占用空间)。
3. 及时释放内存资源
- 任务结果获取后,若无需长期留存内存,可直接写入磁盘存储,再清空内存中的结果列表,避免累积占用内存:
# 示例:将批次结果写入磁盘后清空内存 with open("calc_results.txt", "a") as f: for res in batch_results: f.write(f"{res}\n") batch_results.clear()
- 可通过
ray.internal.free()手动释放不再需要的对象存储引用,加速Ray的垃圾回收。
4. 调整Ray的内存配置
启动Ray时可指定内存相关参数,避免内存过载:
object_store_memory:设置对象存储的最大内存(单位字节),比如ray.init(num_cpus=num_cpu, object_store_memory=10*1024*1024*1024)(10GB)。memory:限制每个Worker进程的最大内存使用,比如ray.init(num_cpus=num_cpu, memory=2*1024*1024*1024)(每个Worker最多2GB)。
内容的提问来源于stack exchange,提问作者Chan
相关产品推荐
相关产品推荐

