【问题标题】:Python PANDAS: Converting from pandas/numpy to dask dataframe/arrayPython PANDAS:从 pandas/numpy 转换为 dask 数据帧/数组
【发布时间】:2018-04-15 02:08:25
【问题描述】:

我正在尝试使用出色的 dask 库将程序转换为可并行化/多线程的。这是我正在转换的程序:

Python PANDAS: Stack by Enumerated Date to Create Records Vectorized

import pandas as pd
import numpy as np
import dask.dataframe as dd
import dask.array as da
from io import StringIO

test_data = '''id,transaction_dt,units,measures
               1,2018-01-01,4,30.5
               1,2018-01-03,4,26.3
               2,2018-01-01,3,12.7
               2,2018-01-03,3,8.8'''

df_test = pd.read_csv(StringIO(test_data), sep=',')
df_test['transaction_dt'] = pd.to_datetime(df_test['transaction_dt'])

df_test = df_test.loc[np.repeat(df_test.index, df_test['units'])]
df_test['transaction_dt'] += pd.to_timedelta(df_test.groupby(level=0).cumcount(), unit='d')
df_test = df_test.reset_index(drop=True)

预期结果:

id,transaction_dt,measures
1,2018-01-01,30.5
1,2018-01-02,30.5
1,2018-01-03,30.5
1,2018-01-04,30.5
1,2018-01-03,26.3
1,2018-01-04,26.3
1,2018-01-05,26.3
1,2018-01-06,26.3
2,2018-01-01,12.7
2,2018-01-02,12.7
2,2018-01-03,12.7
2,2018-01-03,8.8
2,2018-01-04,8.8
2,2018-01-05,8.8 

我突然想到,这可能是尝试并行化的好选择,因为单独的 dask 分区不需要了解彼此的任何信息来完成所需的操作。这是我认为它可能如何工作的天真表示:

dd_test = dd.from_pandas(df_test, npartitions=3)

dd_test = dd_test.loc[da.repeat(dd_test.index, dd_test['units'])]
dd_test['transaction_dt'] += dd_test.to_timedelta(dd.groupby(level=0).cumcount(), unit='d')
dd_test = dd_test.reset_index(drop=True)

到目前为止,我一直在尝试解决以下错误或惯用差异:

  1. “NotImplementedError:仅支持整数值重复。” 我曾尝试将索引转换为 int 列/数组来尝试,但仍然遇到问题。

2. dask 不支持变异操作符:"+="

3. 没有 dask .to_timedelta() 参数

4. 没有 .cumcount() (但我认为 .cumsum() 是可以互换的?!)

如果有任何 dask 专家可以让我知道是否存在阻止我尝试此操作或任何实施技巧的基本障碍,那将是一个很大的帮助!

编辑:

自发布问题以来,我认为我在这方面取得了一些进展:

dd_test = dd.from_pandas(df_test, npartitions=3)
dd_test['helper'] = 1

dd_test = dd_test.loc[da.repeat(dd_test.index, dd_test['units'])]
dd_test['transaction_dt'] = dd_test['transaction_dt'] + (dd.test.groupby('id')['helper'].cumsum()).astype('timedelta64[D]') 
dd_test = dd_test.reset_index(drop=True)

但是,我仍然卡在 dask 数组重复错误上。任何提示仍然欢迎。

【问题讨论】:

  • 什么是df_test['units'].dtype
  • @John Zwinck 它应该是一个整数。
  • 什么是df_test['units'].dtype?请运行代码并查看。不要不检查就说“应该是”。
  • df_test['units'] 和 dd_test['units'] 都是 int64 来确认的。

标签: python pandas numpy dask


【解决方案1】:

不确定这是否正是您要寻找的,但我用 np.repeat 替换了 da.repeat,并将 dd_test.indexdd_test['units'] 显式转换为 numpy 数组,最后将 dd_test['transaction_dt'].astype('M8[us]') 添加到你的 timedelta 计算。

df_test = pd.read_csv(StringIO(test_data), sep=',')

dd_test = dd.from_pandas(df_test, npartitions=3)
dd_test['helper'] = 1

dd_test = dd_test.loc[np.repeat(np.array(dd_test.index), 
np.array(dd_test['units']))]
dd_test['transaction_dt'] = dd_test['transaction_dt'].astype('M8[us]') + (dd_test.groupby('id')['helper'].cumsum()).astype('timedelta64[D]')
dd_test = dd_test.reset_index(drop=True)

df_expected = dd_test.compute()

【讨论】:

  • 哇,感谢您提供有趣的解决方案并且非常有创意。我将研究 dask 和 numpy 如何在此互操作,产生了多少开销,以及这是否充分利用 dask 进行多线程处理。
猜你喜欢
  • 2017-02-04
  • 1970-01-01
  • 2019-08-14
  • 1970-01-01
  • 2018-12-31
  • 2021-04-04
  • 2020-07-28
  • 2015-12-04
  • 2023-03-27
相关资源
最近更新 更多