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

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())
关键修改点
  1. 全局候选集预处理:将name_series_orbis直接转为本地列表name_candidates,避免在延迟函数中引用Dask对象(Dask对象无法在延迟任务中直接序列化执行)。
  2. 按分区处理数据:用dask_df_dealscan.to_delayed()获取每个分区的延迟对象,让每个延迟函数处理整个分区,返回带匹配结果的DataFrame,确保dd.from_delayed能正确拼接成完整的Dask DataFrame。
  3. 返回结构化结果:每个延迟任务返回的是完整的分区DataFrame,而非单个值,这样dd.from_delayed可以识别并组合成统一的DataFrame结构。

内容的提问来源于stack exchange,提问作者Андрей Комиссар</think_never_used_51bce0c785ca2f68081bfa7d91973934>

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 09:33:16