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

调用自定义函数时访问dask_cudf Series属性报错排查

多GPU环境下Dask-CuDF结合Sentence-Transformer编码的AttributeError问题解析

背景

在双T4 GPU的GCP Jupyter Notebook(Debian Python3.8)环境中,尝试用Dask-CuDF的map_partitions实现Sentence-Transformer的并行句子编码,替代之前报错的方案。

核心代码

import cudf
import dask_cudf

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="2GB", rmm_maximum_pool_size="15GB")

client = Client(cluster)

client.run(cudf.set_allocator, "managed") # 使用托管内存替代默认分配器

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]:
        # 报错行:调用to_dask_dataframe()
        emb = sbert_model.encode(sentences=list(df[col].to_dask_dataframe()), batch_size=1024, show_progress_bar=True)
        semb = np.array([str(x) for x in emb])
        emb_array = dask.array.from_array(semb, chunks=chunks)
        df[col+'_emb'] = emb_array
    return df

# 查看数据结构
print(type(dask_df), dask_df.npartitions)
# 输出:(dask_cudf.core.DataFrame, 8)

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'>

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'})

print(new_dask_df.dtypes)
# 输出:
# col_1        object
# col_2        object
# col_1_emb    object
# col_2_emb    object
# dtype: object

# 执行head触发计算时报错
new_dask_df.head()

错误信息

/distributed/client.py:3106: UserWarning: Sending large graph of size 163.19 MiB.
 This may cause some slowdown.
  Consider scattering data ahead of time and using futures.
 warnings.warn(
    2023-06-21 22:43:41,581 - distributed.worker - WARNING - Compute Failed
    Key:       ('test_f_str-61ab17ba9c61c486f2f10e2f92b28b02', 0)
    Function:  subgraph_callable-82156164-6500-4cda-a935-2f8f328a
    args:      (                            col_1                             col_2
    0        This is sentence 6819330  That is another sentence 8466591
    1        This is sentence 7294963  That is another sentence 5338403
    2        This is sentence 7280211  That is another sentence 8222981
    3        This is sentence 2673618  That is another sentence 4661579
    4        This is sentence 6749945  That is another sentence 1511266
    ...                           ...                               ...
    1249995  This is sentence 6905835  That is another sentence 9818177
    1249996   This is sentence 933624  That is another sentence 6111931
    1249997  This is sentence 6352113  That is another sentence 7685898
    1249998  This is sentence 4950656  That is another sentence 4090789
    1249999  This is sentence 3460942  That is another sentence 6747583
[1250000 rows x 2 columns], 'from_cudf-f7c58ae2f4acb6c92d7d79480cb40f46')
    kwargs:    {}
    Exception: 'AttributeError("'Series' object has no attribute 'to_dask_dataframe'")'

矛盾点:单独测试可正常运行

temp = dask_df['col_1'].to_dask_dataframe() # 无报错
list(temp)[:10]
# 输出:
# ['This is sentence 6819330',
#  'This is sentence 7294963',
#  'This is sentence 7280211',
#  'This is sentence 2673618',
#  'This is sentence 6749945',
#  'This is sentence 2843628',
#  'This is sentence 8009669',
#  'This is sentence 391329',
#  'This is sentence 7333531',
#  'This is sentence 6114318']

print(type(temp), type(dask_df['col_1']))
# 输出:(dask.dataframe.core.Series, dask_cudf.core.Series)

问题原因

map_partitions中传入的df参数是CuDF的本地DataFrame分区(而非Dask-CuDF的延迟对象),因此df[col]返回的是cudf.core.Series,而非dask_cudf.core.Series。而cudf.core.Series并没有to_dask_dataframe()方法,这才是报错的根源。

单独测试时,dask_df['col_1']是dask_cudf.core.Series,它确实拥有to_dask_dataframe()方法,因此可以正常执行。但在分区函数内部,处理的是已经计算出的本地CuDF数据,两者的对象类型完全不同。

解决方案

修改分区函数,直接将CuDF Series转换为Python列表,无需通过Dask对象中转:

def test_f_str(df, args):
    col1, col2, chunks = args
    for col in [col1, col2]:
        # 直接将CuDF Series转为列表,供Sentence-Transformer编码
        sentences = df[col].tolist()
        emb = sbert_model.encode(sentences=sentences, batch_size=1024, show_progress_bar=True)
        semb = np.array([str(x) for x in emb])
        # 将数组转为CuDF Series,直接赋值给原DataFrame
        df[col+'_emb'] = cudf.Series(semb)
    return df

额外优化建议

  • 避免在分区函数中创建Dask Array再赋值,直接使用CuDF Series操作更符合GPU数据处理逻辑,减少数据拷贝开销。
  • 模型初始化应放在分区函数内部,通过client.run在每个Worker上单独加载,避免跨进程模型传递的问题:
    def load_model():
        global sbert_model
        sbert_model = SentenceTransformer('all-MiniLM-L6-v2')
    
    client.run(load_model)
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 05:54:58