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

多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内存泄漏的方法

  1. 直接监控硬件指标

    • 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内存泄漏。
  2. 隔离测试验证

    • 单独测试CPU逻辑:临时改用纯Dask CPU集群(替换dask_cuda.LocalCUDACluster为dask.distributed.LocalCluster)运行任务,观察是否仍出现内存泄漏。如果问题消失,说明泄漏点在GPU相关逻辑;如果问题依旧,锁定CPU侧。
    • 单独测试GPU编码逻辑:提取test_f_str中的编码部分,在单GPU环境下循环处理单个分区数据,用nvidia-smi监控内存变化。如果每次循环后GPU内存都增加,可确认是Sentence-Transformers的GPU编码环节存在泄漏。

二、针对当前代码的潜在问题修复

  1. 避免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')
    
  2. 优化内存回收逻辑
    在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
    
  3. 调整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)
    
  4. 减少分区大小
    当前每个分区有125万条数据,单分区处理的batch数量过多,内存占用累积明显。可以把dask_df = dask_cudf.from_cudf(cudf_df, npartitions=8)改成npartitions=16,缩小单分区数据量,降低内存峰值。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 21:38:11