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

Windows系统下使用Dask delayed函数如何解决非托管内存累积问题?

Dask延迟并行化下非托管内存累积问题(Windows环境)

我有一个无数组类型的传统Python函数,需多次运行,因此采用dask.delayed实现Dask并行化,但运行过程中出现非托管内存逐渐累积的情况。

查阅《Tackling unmanaged memory with Dask》博客后,判断该问题由大量Python小对象导致,且进程结束后内存无法自动释放。但由于使用Windows系统,无法测试或应用博客中提到的手动内存修剪方法,在Dask论坛检索相关提问也未找到可行解决方案。

以下是复现该问题的示例代码:

import numpy as np
import gc
import dask
from dask.distributed import Client

client = Client(n_workers=8)

# 需多次调用的原生函数,请勿根据示例的简易性尝试并行化内部逻辑,实际业务代码复杂度更高
def native_function(N):
    curr_sample_1 = frozenset(np.random.default_rng().choice(N, size=500))
    curr_sample_2 = frozenset(np.random.default_rng().choice(N, size=500))
    sum_of_overlap = sum(curr_sample_1.intersection(curr_sample_2))
    
    del curr_sample_1, curr_sample_2  # 避免对象"囤积"导致非托管内存增长
    gc.collect()

    return sum_of_overlap

def another_function(iters):
    list_of_overlap_sums = []
    for i in range(iters):
        native_func_result = dask.delayed(native_function)(10_000)
        list_of_overlap_sums.append(native_func_result)
    
    final_result = dask.delayed(sum)(list_of_overlap_sums).compute()
    
    client.cancel(native_func_result)  # 再次尝试避免非托管内存累积
    client.run(gc.collect)

    return final_result

for i in range(100):
    print(another_function(100))

进程结束时Worker的内存状态如下,可见大部分内存被非托管内存占用,当前虽未超过阈值,但长期运行存在内存溢出隐患:
Worker内存状态

内容的提问来源于stack exchange,提问作者Anirban Chakraborty

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 14:25:15