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的内存状态如下,可见大部分内存被非托管内存占用,当前虽未超过阈值,但长期运行存在内存溢出隐患:
内容的提问来源于stack exchange,提问作者Anirban Chakraborty
相关产品推荐
相关产品推荐

