如何让Dask foldby的每个分组结果存储在独立分区中?
问题:Dask foldby分组结果能否按分组存到独立分区?
我正在尝试理解dask.foldby的工作原理,构建了以下示例:
首先创建一个包含100个元素的Dask Bag:
import dask.bag as db import numpy as np tasks = [{"task_group": i // 10, "numbers": list(range(i, i + 10))} for i in range(1, 1000, 10)] tb = db.from_sequence(tasks)
接着定义聚合函数和合并函数,执行foldby操作计算5个分组的总和:
def aggregate_task(task): return np.array([sum(task['numbers'])] * 10000) # 生成较大的结果数组 def add(x, y): x = aggregate_task(x) if isinstance(x, dict) else x y = aggregate_task(y) if isinstance(y, dict) else y return x + y res = tb.foldby(lambda task: task['task_group'] % 5, add, split_every=25)
计算结果如下:
res.compute() [(0, array([96100, 96100, 96100, ..., 96100, 96100, 96100])), (1, array([98100, 98100, 98100, ..., 98100, 98100, 98100])), (2, array([100100, 100100, 100100, ..., 100100, 100100, 100100])), (3, array([102100, 102100, 102100, ..., 102100, 102100, 102100])), (4, array([104100, 104100, 104100, ..., 104100, 104100, 104100]))]
查看任务图和res.npartitions后发现,所有折叠结果都集中在单个分区中。我的需求是:能否让Dask将每个foldby分组的结果存储在独立分区(对应不同worker)中?这样可以优化内存占用,还能实现并行写入文件。
解决方案
Dask Bag的foldby设计上会将最终的分组聚合结果合并到单个分区,无法直接让每个分组对应独立分区。但可以通过以下两种方式实现需求:
方法1:先聚合再拆分分区
先完成foldby聚合,再将结果拆分为单个分组对应单个分区的Bag:
# 先完成foldby计算 res = tb.foldby(lambda task: task['task_group'] % 5, add, split_every=25) # 将结果拆分为单个元素的分区 res_split = res.map(lambda x: x).repartition(npartitions=5)
此时res_split.npartitions会变为5,每个分区对应一个分组的结果,后续可直接并行执行写入操作。
方法2:改用Dask DataFrame的groupby
Dask DataFrame的分组聚合支持将结果保留为多分区(每个分组一个分区),更适配大体积结果的场景:
import dask.dataframe as dd # 将Bag转为DataFrame df = tb.to_dataframe() # 定义聚合逻辑:先计算每个task的数组,再按分组求和 def agg_func(df): arr = df['numbers'].apply(lambda x: np.array([sum(x)]*10000), meta=('numbers', 'object')) return arr.sum() # 按分组聚合,指定split_out参数等于分组数 res_df = df.groupby(df['task_group'] % 5).agg(agg_func, split_out=5) # 查看分区数 print(res_df.npartitions) # 输出5
split_out=5会让每个分组的结果分配到独立分区,后续可通过to_csv等方法并行写入文件,每个分区对应一个独立文件。
注意事项
- 方法1中,拆分分区的操作是轻量的分区重组,会在计算时自动触发聚合后的拆分逻辑。
- 方法2中,若分组数量不固定,可先通过
df['task_group'] %5.nunique().compute()获取分组数,再设置split_out参数。
内容的提问来源于stack exchange,提问作者rajendra
相关产品推荐
相关产品推荐

