多GPU环境下Dask运行Sentence-Transformers任务的内存泄漏排查
问题描述
在GCP的Jupyter Notebook中,使用2块T4 GPU(单卡15GB)、16个vCPU(60GB内存),通过Sentence-Transformers处理文本数据生成嵌入向量。数据量并不大,但即使已通过Shell设置垃圾回收阈值,Worker节点仍因内存泄漏重启。
代码示例
# 启动dask集群前先在shell执行 export MALLOC_TRIM_THRESHOLD_=65536 !pip install sentence-transformers import os import glob import numpy as np import gc import cudf import dask_cudf import cupy import rmm from dask.distributed import Client, wait, get_worker, get_client from dask_cuda import LocalCUDACluster cluster = LocalCUDACluster(CUDA_VISIBLE_DEVICES="0,1", n_workers=2, threads_per_worker=4, memory_limit="15GB", device_memory_limit="24GB", rmm_pool_size="4GB", rmm_maximum_pool_size="15GB") client = Client(cluster) print(client.run(os.getenv, "MALLOC_TRIM_THRESHOLD_")) # 输出 65536 initial_pool_size = 4*10**9 maximum_pool_size = 15*10**9 rmm.reinitialize(pool_allocator=True, managed_memory=True, initial_pool_size=initial_pool_size, maximum_pool_size=maximum_pool_size, devices=[0,1], logging=True, log_file_name='./tmp/logs/test_sbert_distributed.log') import dask.dataframe as dd import pandas as pd from dask.multiprocessing import get import random df = pd.DataFrame({'col_1': ["This is sentence " + str(x) for x in random.sample(range(10**7), 10**7)], 'col_2': ["That is another sentence " + str(x) for x in random.sample(range(10**7), 10**7)]}) cudf_df = cudf.DataFrame.from_pandas(df) dask_df = dask_cudf.from_cudf(cudf_df, npartitions=8) from sentence_transformers import SentenceTransformer import numpy as np sbert_model = SentenceTransformer('all-MiniLM-L6-v2') def test_f_str(df, args): col1, col2, chunks = args for col in [col1, col2]: emb = sbert_model.encode(sentences=df[col].to_arrow().to_pylist(), batch_size=1250, show_progress_bar=True) semb = np.array([str(x) for x in emb]) df[col+'_emb'] = semb return df dask_cudf.core.Series chunks = dask_df.map_partitions(lambda x: len(x)).compute().to_numpy() print(chunks, type(chunks)) # 输出:[1250000 1250000 1250000 1250000 1250000 1250000 1250000 1250000] <class 'numpy.ndarray'> dask_df.npartitions, dask_df.persist() # 输出:(8, <dask_cudf.DataFrame | 8 tasks | 8 npartitions>) new_dask_df = dask_df.map_partitions(test_f_str, args=('col_1', 'col_2', chunks), meta={'col_1':'object', 'col_2':'object', 'col_1_emb':'object', 'col_2_emb':'object'}) new_dask_df.dtypes # 输出: # col_1 object # col_2 object # col_1_emb object # col_2_emb object # dtype: object new_dask_df.compute() # 报错:WARNING - Unmanaged memory use is high. This may indicate a memory leak or the memory may not be released to the OS # -- Unmanaged memory: 9.96 GiB -- Worker memory limit: 13.97 GiB
已尝试的无效方案
- 调整Dask非托管内存相关配置与优化手段
- 排查Dask内存回收机制问题
- 尝试Dask内存泄漏的常见临时解决方法
更新内容
目前使用Jupyter Lab中的NVIDIA GPU仪表盘监控,Worker内存图表未显示未管理或泄漏的内存,但GPU内存图表变为橙色并提示内存溢出。
请问如何确认这是CPU还是GPU内存泄漏?
解决方案
一、区分CPU/GPU内存泄漏的方法
直接监控硬件指标
- CPU内存:在终端执行
htop或free -h,实时查看每个Worker进程的常驻内存(RSS)占用,对比Dask仪表盘的Worker内存数据。如果系统层面CPU内存持续上涨且任务结束后不回落,结合Dask的非托管内存警告,可判定为CPU内存泄漏。 - GPU内存:执行
nvidia-smi -l 1(每秒刷新一次),观察每张T4 GPU的Used GPU Memory列。如果任务执行过程中GPU内存持续攀升,即使任务完成后也未释放,或超过设置的device_memory_limit,则属于GPU内存泄漏。
- CPU内存:在终端执行
隔离测试验证
- 单独测试CPU逻辑:临时改用纯Dask CPU集群(替换
dask_cuda.LocalCUDACluster为dask.distributed.LocalCluster)运行任务,观察是否仍出现内存泄漏。如果问题消失,说明泄漏点在GPU相关逻辑;如果问题依旧,锁定CPU侧。 - 单独测试GPU编码逻辑:提取
test_f_str中的编码部分,在单GPU环境下循环处理单个分区数据,用nvidia-smi监控内存变化。如果每次循环后GPU内存都增加,可确认是Sentence-Transformers的GPU编码环节存在泄漏。
- 单独测试CPU逻辑:临时改用纯Dask CPU集群(替换
二、针对当前代码的潜在问题修复
避免CPU-GPU数据频繁转换
当前代码中df[col].to_arrow().to_pylist()会把GPU上的cudf数据转成CPU端的Python列表,不仅效率低,还会在CPU侧产生大量临时对象,增加内存压力。建议优化数据传递方式:# 直接传递cudf转成的pandas数据,避免转成Python列表 emb = sbert_model.encode(sentences=df[col].to_pandas(), batch_size=1250, show_progress_bar=True, device='cuda')优化内存回收逻辑
在test_f_str函数末尾手动清理临时变量,并触发垃圾回收:def test_f_str(df, args): col1, col2, chunks = args for col in [col1, col2]: emb = sbert_model.encode(sentences=df[col].to_arrow().to_pylist(), batch_size=1250, show_progress_bar=True) semb = np.array([str(x) for x in emb]) df[col+'_emb'] = semb # 清理临时变量并触发GC del emb, semb gc.collect() return df调整RMM配置
当前rmm.reinitialize中设置了managed_memory=True,可能导致GPU内存管理混乱。建议关闭托管内存,改用纯池化分配:rmm.reinitialize(pool_allocator=True, managed_memory=False, initial_pool_size=initial_pool_size, maximum_pool_size=maximum_pool_size, devices=[0,1], logging=True)减少分区大小
当前每个分区有125万条数据,单分区处理的batch数量过多,内存占用累积明显。可以把dask_df = dask_cudf.from_cudf(cudf_df, npartitions=8)改成npartitions=16,缩小单分区数据量,降低内存峰值。
内容的提问来源于stack exchange,提问作者mtnt
相关产品推荐
相关产品推荐

