You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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. 这种磁盘占用异常的情况是否属于正常行为?
  2. 该如何解决这个问题?
  3. 强制垃圾回收是否可行?

解答

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.19 17:50:25