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

如何让Dask map_partitions返回多行DataFrame?含Delayed方案需求

解决Dask/Pandas中应用函数生成多行DataFrame的问题

先搞定Pandas里的问题

假设你的自定义函数是这样的(输入单个值,返回2行DataFrame):

import pandas as pd

def cell_func(val):
    return pd.DataFrame({'new_col1': [val+1, val+2], 'new_col2': [val*2, val*3]})

错误用法(只会得到列名行)

直接用applymap或者简单apply的话,每个单元格会存一个DataFrame对象,根本不会展开成行:

df = pd.DataFrame({'col1': [10, 20], 'col2': [30, 40]})
# 这么写只会得到每个单元格是DataFrame的结果,不是你要的多行数据
bad_result = df.applymap(cell_func)

正确操作:扁平化展开结果

要把每个单元格生成的DataFrame拆成行,再合并成大表:

# 1. 把原表转成长格式,每个单元格对应一行,保留原行列索引
stacked = df.stack().reset_index(name='original_val')
# 2. 对每个原始值应用函数,生成多行结果
stacked['result'] = stacked['original_val'].apply(cell_func)
# 3. 展开result列里的DataFrame,再和原索引列合并
final_df = stacked.explode('result').reset_index(drop=True)
final_df = pd.concat([final_df.drop('result', axis=1), final_df['result'].apply(pd.Series)], axis=1)

这样就能得到每个单元格生成的多行数据,还能保留原始的行列位置信息(如果需要的话)。


Dask方案:优先用Dask Delayed实现

你要求用Delayed方案,确实它更灵活,适合这种单元格级展开的场景。

步骤1:定义处理函数和加载数据

import dask
import dask.dataframe as dd
from dask.delayed import delayed

# 复用之前的cell_func
def cell_func(val):
    return pd.DataFrame({'new_col1': [val+1, val+2], 'new_col2': [val*2, val*3]})

# 定义处理单个分区的函数(和Pandas里的逻辑一致)
def process_partition(partition):
    stacked = partition.stack().reset_index(name='original_val')
    stacked['result'] = stacked['original_val'].apply(cell_func)
    final_part = stacked.explode('result').reset_index(drop=True)
    return pd.concat([final_part.drop('result', axis=1), final_part['result'].apply(pd.Series)], axis=1)

# 示例Dask DataFrame(也可以从文件加载)
dask_df = dd.from_pandas(pd.DataFrame({'col1': [10,20,30,40], 'col2': [50,60,70,80]}), npartitions=2)

步骤2:用Delayed处理分区并组装结果

# 把每个分区转为延迟任务,应用处理函数
delayed_partitions = [delayed(process_partition)(part) for part in dask_df.to_delayed()]

# 将延迟的分区结果转为Dask DataFrame
final_dask_df = dd.from_delayed(delayed_partitions)

# 触发计算,得到最终结果
final_result = final_dask_df.compute()

为什么Delayed更适合?

  • 完全自定义每个分区的处理流程,避免map_partitions可能出现的结果对齐bug
  • 能更精细地控制任务的依赖关系和并行策略
  • 面对复杂的单元格级转换时,逻辑比map_partitions更清晰

Dask map_partitions的正确用法(备选)

如果一定要用map_partitions,核心逻辑和Delayed一样,确保处理函数返回结构一致的DataFrame:

def process_partition_map(partition):
    stacked = partition.stack().reset_index(name='original_val')
    stacked['result'] = stacked['original_val'].apply(cell_func)
    final_part = stacked.explode('result').reset_index(drop=True)
    return pd.concat([final_part.drop('result', axis=1), final_part['result'].apply(pd.Series)], axis=1)

# 应用map_partitions并计算
result_dask = dask_df.map_partitions(process_partition_map).compute()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 06:40:26