【问题标题】:Python Dask map_partitionsPython Dask map_partitions
【发布时间】:2019-01-07 04:37:53
【问题描述】:

可能是这个question 的延续,从map_partitions 的dask 文档示例开始工作。

import dask.dataframe as dd
df = pd.DataFrame({'x': [1, 2, 3, 4, 5],     'y': [1., 2., 3., 4., 5.]})
ddf = dd.from_pandas(df, npartitions=2)

from random import randint

def myadd(df):
    new_value = df.x + randint(1,4)
    return new_value

res = ddf.map_partitions(lambda df: df.assign(z=myadd)).compute()
res

在上面的代码中,randint 只被调用一次,而不是像我期望的那样每行调用一次。怎么会?

输出:

X Y Z

1 1 4

2 2 5

3 3 6

4 4 7

5 5 8

【问题讨论】:

    标签: python pandas dask


    【解决方案1】:

    如果您在原始 pandas 数据帧上执行相同的操作 (df.x + randint(1,4)),您只会得到一个随机数,添加到该列的每个先前值。这与 pandas 的情况完全相同,只是它为每个分区调用一次 - 这就是 map_partition 所做的。

    如果您想为每一行添加一个新的随机数,您应该首先考虑如何使用 pandas 来实现这一点。我马上能想到两个:

    df.x.map(lambda x: x + random.randint(1, 4))
    

    df.x + np.random.randint(1, 4, size=len(df.x))
    

    如果您将 newvalue = 行替换为其中之一,它将按预期工作。

    【讨论】:

      猜你喜欢
      • 2021-09-30
      • 1970-01-01
      • 1970-01-01
      • 2022-08-06
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-07-02
      • 2022-01-26
      相关资源
      最近更新 更多