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

如何让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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.23 07:54:18