调用自定义函数时访问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
相关产品推荐
相关产品推荐

