如何高效计算Dask DataFrame各分区的逐元素平均值?
问题描述
我有一个由多个二维DataFrame组成的Dask DataFrame(每个分区为一个二维DataFrame),希望找到最高效的方法计算所有分区的逐元素平均值,得到单个二维DataFrame。示例如下:
partition 1 partition 2 partition 3 A B A B A B 0 | 1 4 0 | 7 10 0 | 13 16 1 | 2 5 1 | 8 11 1 | 14 17 2 | 3 6 2 | 9 12 2 | 15 18
需要对所有分区的对应位置元素求平均,例如:
- 第0行A列:mean(1,7,13) = 7
- 第1行A列:mean(2,8,14) = 8
- 第0行B列:mean(4,10,16) = 10
- 以此类推
最终得到的结果为:
A B 0 | 7.0 10.0 1 | 8.0 11.0 2 | 9.0 12.0
注:示例中原Row 1 Column A的计算式笔误已修正。我尝试过使用dask.sum()后除以分区数,但该方法仅支持按行或列求和;map_partition函数也不适用,因为该操作并非独立执行。
解决方案
针对这种每个分区结构完全一致(行索引、列名完全对应)的场景,推荐以下两种高效实现方式:
方法一:按索引分组求平均
利用groupby对行索引分组,直接计算每组的平均值,这是最直观且易维护的方式:
import dask.dataframe as dd import pandas as pd # 构造示例Dask DataFrame df1 = pd.DataFrame({'A': [1,2,3], 'B': [4,5,6]}) df2 = pd.DataFrame({'A': [7,8,9], 'B': [10,11,12]}) df3 = pd.DataFrame({'A': [13,14,15], 'B': [16,17,18]}) ddf = dd.from_pandas(df1, npartitions=1) ddf = ddf.append(df2).append(df3) # 按行索引分组计算逐元素平均 result = ddf.groupby(ddf.index).mean().compute() print(result)
方法二:先求和再除以分区数(更高效)
如果分区数量固定且已知,可以先计算所有分区对应位置的元素和,再除以分区数,减少分组操作的开销:
import dask.dataframe as dd import pandas as pd # 构造示例Dask DataFrame df1 = pd.DataFrame({'A': [1,2,3], 'B': [4,5,6]}) df2 = pd.DataFrame({'A': [7,8,9], 'B': [10,11,12]}) df3 = pd.DataFrame({'A': [13,14,15], 'B': [16,17,18]}) ddf = dd.from_pandas(df1, npartitions=1) ddf = ddf.append(df2).append(df3) # 获取分区数量 n_partitions = ddf.npartitions # 计算所有分区的逐元素和,再除以分区数得到平均 sum_df = ddf.map_partitions(lambda df: df).sum(axis=0, level=0).compute() result = sum_df / n_partitions print(result)
说明
- 两种方法的核心都是利用每个分区的行索引对应关系,确保同一位置的元素被聚合计算。
- 方法二更适合数据量较大的场景,因为求和操作的计算开销比分组平均更低。
内容的提问来源于stack exchange,提问作者solidahmad
相关产品推荐
相关产品推荐

