Ray对象溢出累积问题咨询:是否正常及优化方案
Ray并行计算磁盘溢出异常问题解答
问题背景
使用Ray进行计算并行化时,处理约1.6 GiB的模拟数据,磁盘占用却超过130 GiB,无法接受。运行日志、代码及环境如下:
运行日志
Number of steps: 55 (9,091 simulations each) 0% [2m[36m(raylet)[0m Spilled 3702 MiB, 12 objects, write throughput 661 MiB/s. Set RAY_verbose_spill_logs=0 to disable this message. [2m[36m(raylet)[0m Spilled 5542 MiB, 17 objects, write throughput 737 MiB/s. 2% [2m[36m(raylet)[0m Spilled 9883 MiB, 33 objects, write throughput 849 MiB/s. 5% [2m[36m(raylet)[0m Spilled 16704 MiB, 58 objects, write throughput 997 MiB/s. 13% [2m[36m(raylet)[0m Spilled 32903 MiB, 124 objects, write throughput 784 MiB/s. 29% [2m[36m(raylet)[0m Spilled 66027 MiB, 268 objects, write throughput 661 MiB/s. 53% [2m[36m(raylet)[0m Spilled 131920 MiB, 524 objects, write throughput 461 MiB/s. 60%
运行代码
def get_res_parallel(simulations, num_loads=num_cpus): load_size = simulations.shape[0] / num_loads simulations_per_load = [simulations[round(n * load_size): round((n+1) * load_size)] for n in range(num_loads)] # 2D numpy arrays results = ray.get([get_res_ray.remote(simulations=simulations) for simulations in simulations_per_load]) return np.vstack(results) MAX_RAM = 6 * 2**30 # 6 GiB def get_expected_res(simulations, MAX_RAM=MAX_RAM): expected_result = np.zeros(shape=87_381, dtype=np.float64) bytes_per_res = len(expected_result) * (64 // 8) num_steps = simulations.shape[0] * bytes_per_res // MAX_RAM + 1 step_size = simulations.shape[0] / num_steps print(f"Number of steps: {num_steps} ({step_size:,.0f} simulations each)") for n in range(num_steps): print(f"\r{n / num_steps:.0%}", end="") step_simulations = simulations[round(n * step_size): round((n+1) * step_size)] results = get_res_parallel(simulations=step_simulations) expected_result += results.mean(axis=0) print(f"\r100%") return expected_result / num_steps
运行环境
Mac M1(16 GiB内存)、Ray 2.0.0、Python 3.9.13
疑问
- 这种磁盘占用异常的情况是否属于正常行为?
- 该如何解决这个问题?
- 强制垃圾回收是否可行?
解答
1. 是否属于正常行为?
这绝对不属于正常行为。1.6 GiB输入数据产生130 GiB溢出数据,说明Ray对象存储中存在大量未及时清理的对象,大概率是代码逻辑导致的引用泄漏,或是Ray 2.0.0版本的溢出机制缺陷。
2. 解决方案
(1)显式释放Ray对象与本地引用
Ray默认不会自动删除对象存储中的远程任务结果,需手动释放;同时本地大数组也需及时清理:
- 修改
get_res_parallel,调用ray.internal.free释放远程任务引用:
def get_res_parallel(simulations, num_loads=num_cpus): load_size = simulations.shape[0] / num_loads simulations_per_load = [simulations[round(n * load_size): round((n+1) * load_size)] for n in range(num_loads)] # 保存远程任务引用 refs = [get_res_ray.remote(simulations=s) for s in simulations_per_load] results = ray.get(refs) # 释放Ray对象存储中的数据 ray.internal.free(refs) return np.vstack(results)
- 在
get_expected_res循环中,计算完成后立即删除本地数组并触发GC:
for n in range(num_steps): print(f"\r{n / num_steps:.0%}", end="") step_simulations = simulations[round(n * step_size): round((n+1) * step_size)] results = get_res_parallel(simulations=step_simulations) expected_result += results.mean(axis=0) # 清理本地大数组并触发GC del results import gc gc.collect()
(2)升级Ray版本
Ray 2.0.0是较旧版本,后续版本(如2.5+)对溢出机制和对象GC做了大量优化,修复了诸多内存/磁盘泄漏问题。执行以下命令升级:
pip install --upgrade ray
(3)调整Ray资源配置
通过初始化参数或环境变量限制内存与磁盘溢出:
- 初始化Ray时设置内存与对象存储大小:
ray.init(memory=4*1024**3, object_store_memory=2*1024**3) # 4GiB进程内存,2GiB对象存储
- 设置环境变量限制最大溢出磁盘空间:
export RAY_spill_max_size=20*1024**3 # 限制溢出至20GiB
(4)优化计算逻辑,减少中间对象
当前代码返回完整结果数组再求均值,可改为远程任务直接返回均值,大幅降低对象大小:
- 修改远程任务
get_res_ray:
@ray.remote def get_res_ray(simulations): # 原逻辑生成full_results后,直接返回均值 full_results = ... # 原有计算逻辑 return full_results.mean(axis=0)
- 修改
get_res_parallel直接聚合均值:
def get_res_parallel(simulations, num_loads=num_cpus): load_size = simulations.shape[0] / num_loads simulations_per_load = [simulations[round(n * load_size): round((n+1) * load_size)] for n in range(num_loads)] refs = [get_res_ray.remote(simulations=s) for s in simulations_per_load] step_means = ray.get(refs) ray.internal.free(refs) # 对子任务均值再求均值,得到当前step的总体均值 return np.mean(step_means, axis=0)
- 调整
get_expected_res中的循环逻辑:
for n in range(num_steps): print(f"\r{n / num_steps:.0%}", end="") step_simulations = simulations[round(n * step_size): round((n+1) * step_size)] step_mean = get_res_parallel(simulations=step_simulations) expected_result += step_mean del step_mean gc.collect()
此优化后,每个远程任务仅返回约690KB的数组,几乎不会产生磁盘溢出。
3. 强制垃圾回收是否可行?
强制Python GC有一定作用,但仅能清理本地进程中的对象引用,无法直接清理Ray对象存储中的数据。必须配合ray.internal.free释放Ray管理的远程对象,否则对象存储数据会持续占用磁盘。Ray自身有GC机制,但旧版本触发不及时,显式调用ray.internal.free更可靠。
内容的提问来源于stack exchange,提问作者Louis-Amand
相关产品推荐
相关产品推荐

