【问题标题】:How to return one NumPy array per partition in Dask?如何在 Dask 中为每个分区返回一个 NumPy 数组?
【发布时间】:2021-03-04 12:10:01
【问题描述】:

我需要计算许多 NumPy 数组(最多可以是 4 维),一个用于 Dask 数据帧的每个分区,然后将它们添加为数组。但是,我正在努力让 map_partitions 为每个分区返回一个数组,而不是为所有分区返回一个数组。

import dask.dataframe as dd
import numpy as np, pandas as pd

df = pd.DataFrame(range(15), columns=['x'])
ddf = dd.from_pandas(df, npartitions=3)

def func(partition):
    # Here I also tried returning the array in a list and in a tuple
    return np.array([[1, 2], [3, 4]])

# Here I tried all the options available for 'meta'
results = ddf.map_partitions(func).compute()

那么results就是:

array([[1, 2],
       [3, 4],
       [1, 2],
       [3, 4],
       [1, 2],
       [3, 4]])

如果我改为results.sum().compute(),我会得到30

我想得到的是:

[np.array([[1, 2],[3, 4]]), np.array([[1, 2],[3, 4]]), np.array([[1, 2],[3, 4]])]

所以如果我计算总和,我会得到:

array([[ 3,  6],
       [ 9, 12]])

您如何使用 Dask 实现此结果?

【问题讨论】:

    标签: dataframe numpy dask


    【解决方案1】:

    我设法让它像这样工作,但我不知道这是否是最好的方法:

    from dask import delayed
    results = []
    for partition in ddf.partitions:
        result = delayed(func)(partition)
        results.append(result)
    
    delayed(sum)(results).compute()
    

    计算结果为:

    array([[ 3,  6],
           [ 9, 12]])
    

    【讨论】:

      【解决方案2】:

      你是对的,一个 dask-array 通常被视为一个单独的逻辑数组,它恰好是由片段组成的。如果您没有使用逻辑层,您可以单独使用delayed 完成您的工作。另一方面,您想要的最终结果似乎是所有数据的总和,所以更简单的可能是合适的reshapesum(axis=)

      ddf.map_partitions(func).compute_chunk_sizes().reshape(
          -1, 2, 2).sum(axis=0).compute()
      

      (需要compute_chunk_sizes,因为虽然您的原始 pandas 数据框的大小已知,但 Dask 尚未评估您的函数,还不知道它返回的大小)

      但是,根据您的设置,以下内容将起作用并且与您最初的尝试更相似,请参阅.to_delayed()

      list_of_delayed = ddf.map_partitions(func).to_delayed().tolist()
      tuple_of_np_lists = dask.compute(*list_of_delayed)
      

      tolist 强制评估包含的延迟对象)

      【讨论】:

      • 嗨,我试过你的解决方案,但似乎都只适用于二维数组。在我的例子中,数组最多可以是 4 维的。两者都有我得到KeyError
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-06-10
      • 1970-01-01
      • 1970-01-01
      • 2015-07-02
      • 2022-01-08
      相关资源
      最近更新 更多