【问题标题】:Dask dataframe join slow as pandasDask 数据框加入慢如熊猫
【发布时间】:2019-07-08 13:19:56
【问题描述】:

我有 2 个数据帧,一个称为animes ~10k 行数据,一个称为animelists ~30M 行数据,并且想要加入它们。我用 pandas 对它进行了基准测试,它的速度只有 7% 左右,这并不多,我想知道如果我有 16 个核心是否会更快。

我有 pandas 数据框,我在其中设置索引

animes = animes.set_index('anime_id')
animelists = animelists.set_index('anime_id')

数据看起来是这样的(我省略了其他列),动漫:

anime_id | genres
-------- | ------
11013    | Comedy, Supernatural, Romance, Shounen
2104     | Comedy, Parody, Romance, School, Shounen
5262     | Comedy, Magic, School, Shoujo

和动漫爱好者:

anime_id | username | my_score
21       | karthiga | 9
59       | karthiga | 7
74       | karthiga | 7

然后我由此创建了 Dask Dataframes

animes_dd = dd.from_pandas(animes, npartitions=8)
animelists_dd = dd.from_pandas(animelists, npartitions=8)

我想将各个动漫流派与动漫列表有效地结合起来,以按流派查询分数。我在 pandas 中有代码可以做到这一点:

genres_arr = animes['genres'].str.replace(' ', '').str.split(',', expand=True).stack().reset_index(drop=True, level=1).to_frame(name='genre')
genres_arr = genres_arr[genres_arr['genre'] != '']
resulting_df = animelists.merge(genres_arr, how='inner', left_index=True, right_index=True)
# this takes 1min 37s

和 dask 中的相同代码:

genres_arr_dd = animes_dd['genres'].map_partitions(lambda x: x.str.replace(' ', '').str.split(',', expand=True).stack().reset_index(drop=True, level=1)).to_frame(name='genre')
genres_arr_dd = genres_arr_dd[genres_arr_dd['genre'] != '']
resulting_dd = animelists_dd.merge(genres_arr_dd, how='inner', left_index=True, right_index=True).compute()
# this takes 1min 30s

(生成的数据帧有大约 1.4 亿行)

有什么方法可以加快速度吗?我遵循official performance guide,对索引列执行连接,每个 Dask Dataframe 上有 8 个分区,因此应该为有效的多处理连接做好准备。

这里出了什么问题,我应该如何加快速度?

当我在 jupyter notebook 中运行代码时,我正在观察每个核心的 CPU 使用率,它非常低,有时只有一个核心处于活动状态,并且以 100% 的速度运行。好像并行不好。

【问题讨论】:

  • 你得到了加速,万岁!
  • 是的,但可以忽略不计,我认为多线程与单线程版本相比,在应该很好地并行化的任务上显示出更大的加速。如果这是 Dask 的最大能力,那么 Dask 就让人大失所望。
  • 并行不是魔法,有大量的取舍
  • 我知道并行不是魔法。仅使用具有 8 个线程的 pandas 和 joblib,我设法优化了代码,因此我在 1 分 4 秒内获得了相同的数据集,这更显着加速。就像我对 Dask 所期望的那样,显然直接使用多线程来执行此类任务会更有效。
  • 不,多进程,这是重点

标签: python multithreading pandas dataframe dask


【解决方案1】:

这在其他地方已经重复了,所以我会保持很简短。

  • from_pandas->compute 表示您正在往返所有数据;你想加载工人(例如,dd.read_csv)并在工人中聚合,而不是移动整个数据集

  • 调度程序的选择非常重要。如果您的系统监视器说您正在使用一个 CPU,那么您可能受到 GIL 的限制,应该尝试使用分布式调度程序,并使用适当的进程/线程组合。它还将在其仪表板上为您提供有关正在发生的事情的更多诊断信息

  • Pandas 速度很快,当数据量较小时,dask 的额外开销虽然也很小,但可能超过您获得的任何并行性。

【讨论】:

  • 哦,我明白了。我不知道 from_pandas->compute 是问题,我从未在文档中看到我不应该使用它。我认为它将数据拆分为 8 个工作人员。谢谢,我将尝试使用 dask 加载数据,并稍微修改一下调度程序。我知道 Pandas 在小数据上的速度很快,而且我实际上观察到,在 pandas 中对约 10-50k 行的操作比在 dask 中更快。我认为 30M 到 140M 行应该足够大,可以看出差异。
  • 所以我尝试了所有可能性,Pandas + Joblib 做得最快。
  • dask 与进程分布基本相同。多处理调度程序与 joblib 更相似,但不再经常使用,因为分布式调度程序为您提供了所有不错的诊断。
猜你喜欢
  • 2020-11-24
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-08-10
  • 2013-12-20
相关资源
最近更新 更多