Dask from_delayed返回标量而非DataFrame的问题求助
问题分析
你遇到的核心问题有两个:
- 每个
delayed函数返回的是单个字符串值,而非DataFrame/Series对象,dd.from_delayed需要每个延迟对象返回可组合成DataFrame的结构(比如单行DataFrame、Series),否则会将所有标量合并成一个Dask标量。 - 错误地将Dask Series转为延迟列表后逐行遍历,违背了Dask的分区并行逻辑,这种方式不仅效率低下,还会导致分区信息丢失。
修正方案
我们需要调整代码,让每个延迟任务处理整个分区的数据,返回可拼接的DataFrame/Series,同时确保候选名称集能被每个任务访问到:
import dask.dataframe as dd import pandas as pd from fuzzywuzzy import process from dask import delayed # 第一个数据集 data_orbis = {'Name': ['Johanna', 'Chanti', 'Pamela', 'Connie', 'Cesar', 'Eldar', 'Chris', "Richard","Rechard","Carl","Carl"], 'ID': ['100', "200","300","400","500","600","700", "800", "900", "999", "1100" ], 'Country': ["Germany", "Germany", "Czechia","Russia","Andorra","France","Germany","Germany","Germany", "Denmark", "Sweden" ], 'City': ["Mainz", "Mainz", "Brno", "Moscow", "Andorra","Paris", "Berlin", "Mainz", "Mannheim","Copenhagen", "Malmo" ]} df_orbis = pd.DataFrame(data_orbis) # 因为需要全局匹配,直接将候选名称转为本地列表(大数据集可考虑广播) name_candidates = df_orbis["Name"].tolist() # 第二个数据集 data_dealscan = {'Name': ['Andrey', 'Canti', 'Pamelo', 'Cannie', 'Cezar', 'Eltor', 'Chriss', "Richard", "Rechard","Carl","Carll"], 'ID': ['np.nan', "np.nan","np.nan","np.nan","500","600","700", "800","900", "999", "np.nan"], 'Country': ["Germany", "Germany", "Czechia","Russia","Andorra","France","Germany","Germany","Germany","Denmark", "Sweden"], 'City': ["Mainz", "Mainz", "Brno", "Moscow", "Andorra","Paris", "Berlin", "Mainz", "Mannheim","Copenhagen", "Malmo" ]} df_dealscan = pd.DataFrame(data_dealscan) dask_df_dealscan = dd.from_pandas(df_dealscan, npartitions=4) # 定义处理整个分区的延迟函数 @delayed def match_partition(partition_df, candidates): # 对分区内每一行执行模糊匹配 partition_df['Matched_Name'] = partition_df['Name'].apply( lambda x: process.extractOne(x, candidates)[0] ) return partition_df # 对每个分区应用延迟函数 lazy_partitions = [] for part in dask_df_dealscan.to_delayed(): lazy_part = match_partition(part, name_candidates) lazy_partitions.append(lazy_part) # 将延迟的分区结果组合成Dask DataFrame ddf = dd.from_delayed(lazy_partitions) # 验证结果 print(ddf.compute())
关键修改点
- 全局候选集预处理:将
name_series_orbis直接转为本地列表name_candidates,避免在延迟函数中引用Dask对象(Dask对象无法在延迟任务中直接序列化执行)。 - 按分区处理数据:用
dask_df_dealscan.to_delayed()获取每个分区的延迟对象,让每个延迟函数处理整个分区,返回带匹配结果的DataFrame,确保dd.from_delayed能正确拼接成完整的Dask DataFrame。 - 返回结构化结果:每个延迟任务返回的是完整的分区DataFrame,而非单个值,这样
dd.from_delayed可以识别并组合成统一的DataFrame结构。
内容的提问来源于stack exchange,提问作者Андрей Комиссар</think_never_used_51bce0c785ca2f68081bfa7d91973934>
相关产品推荐
相关产品推荐

