【问题标题】:How to execute a multi-threaded `merge()` with dask? How to use multiples cores via qsub?如何使用 dask 执行多线程`merger()`?如何通过 qsub 使用多核?
【发布时间】:2017-02-24 11:10:11
【问题描述】:

我刚刚开始使用 dask,但我仍然对如何使用多线程或使用集群执行简单的 pandas 任务感到困惑。

让我们以pandas.merge()dask 数据帧为例。

import dask.dataframe as dd

df1 = dd.read_csv("file1.csv")
df2 = dd.read_csv("file2.csv")

df3 = dd.merge(df1, df2)

现在,假设我要在具有 4 个内核的笔记本电脑上运行它。如何为这个任务分配 4 个线程?

看来正确的做法是:

dask.set_options(get=dask.threaded.get)
df3 = dd.merge(df1, df2).compute()

这将使用尽可能多的线程(即,笔记本电脑上有尽可能多的具有共享内存的内核,4)?如何设置线程数?

假设我在一个拥有 100 个内核的设施中。如何以与使用qsub 向集群提交作业相同的方式提交此文件? (类似于通过 MPI 在集群上运行任务?)

dask.set_options(get=dask.threaded.get)
df3 = dd.merge(df1, df2).compute

【问题讨论】:

  • 你试过dd.merge(df1, df2).compute(num_workers=4) 吗?
  • @Boud 您在文档中的什么地方遇到过这个?
  • @Boud 谢谢。我会删除这个问题---虽然我仍然不确定 MPI 方面:)
  • 不要删除它可能对其他人有用

标签: python multithreading pandas cluster-computing dask


【解决方案1】:

单机调度

默认情况下,Dask.dataframe 将使用线程调度程序,其线程数与您机器中的逻辑内核数一样多。

正如 cmets 中所指出的,您可以使用 .compute() 方法的关键字参数来控制线程数或 Pool 实现。

分布式机器调度

您可以使用dask.distributeddeploy dask workers across many nodes in a cluster。使用qsub 执行此操作的一种方法是在本地启动dask-scheduler

$ dask-scheduler
Scheduler started at 192.168.1.100:8786

然后使用qsub 启动多个dask-worker 进程,指向报告的地址:

$ qsub dask-worker 192.168.1.100:8786 ... <various options>

截至昨天,有一个实验包可以在任何支持 DRMAA 的系统(包括 SGE/qsub-like 系统)上执行此操作:https://github.com/dask/dask-drmaa

完成此操作后,您可以创建一个dask.distributed.Client 对象,它将接管作为默认调度程序

from dask.distributed import Client
c = Client('192.168.1.100:8786')  # now computations run by default on the cluster

多线程性能

请注意,从 Pandas 0.19 版开始,pd.merge 的 GIL 仍未发布,因此我不希望使用多线程来大幅提升速度。如果这对您很重要,那么我建议您在此处发表评论:https://github.com/pandas-dev/pandas/issues/13745

【讨论】:

  • 感谢您的帮助! “所以我不希望使用多线程来大幅提升速度”我主要使用 dash 来解决 RAM 问题。虽然加速merge 会很棒,但目前我只是想完成一个简单的merge 任务。
猜你喜欢
  • 1970-01-01
  • 2020-10-15
  • 2018-05-17
  • 1970-01-01
  • 2020-12-14
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2011-08-03
相关资源
最近更新 更多