【问题标题】:Analyzing data flow of Dask dataframes分析 Dask 数据帧的数据流
【发布时间】:2018-10-12 05:21:07
【问题描述】:

我有一个数据集存储在一个制表符分隔的文本文件中。该文件如下所示:

date    time    temperature
2010-01-01  12:00:00    10.0000 
...

temperature 列包含以摄氏度 (°C) 为单位的值。 我使用 Dask 计算每日平均温度。这是我的代码:

from dask.distributed import Client
import dask.dataframe as dd

client = Client("<scheduler URL")
inputDataFrame = dd.read_table("<input file>").drop('time', axis=1)
groupedData = inputDataFrame.groupby('date')
meanDataframe = groupedData.mean()
result = meanDataframe.compute()
result.to_csv('result.out', sep='\t')

client.close()

为了提高我程序的性能,想了解一下Dask数据帧造成的数据流。

  1. read_table()如何将文本文件读入数据框?客户端是否读取整个文本文件并将数据发送到调度程序,调度程序将数据分区并将其发送给工作人员?还是每个工作人员都直接从文本文件中读取其工作的数据分区?
  2. 在创建中间数据帧时(例如,通过调用drop())是否会将整个中间数据帧发送回客户端,然后发送给工作人员进行进一步处理?
  3. 组的相同问题:组对象的数据在哪里创建和存储?它如何在客户端、调度程序和工作人员之间流动?

我提出问题的原因是,如果我使用 Pandas 运行类似的程序,计算速度大约快两倍,我试图了解导致 Dask 开销的原因。由于结果数据帧的大小与输入数据的大小相比非常小,我认为在客户端、调度程序和工作人员之间移动输入和中间数据会产生相当多的开销。

【问题讨论】:

    标签: pandas dask dask-distributed


    【解决方案1】:

    1) 数据由工人读取。客户端确实提前阅读了一些内容,以找出列名和类型,并可选地找到用于拆分文件的行分隔符。请注意,所有工作人员都必须能够访问感兴趣的文件,这在集群上工作时可能需要一些共享文件系统。

    2), 3) 实际上,dropgroupbymean 方法根本不会生成中间数据帧,它们只是累积要执行的操作图(即,它们是 懒惰)。你可以为这些步骤计时,看看它们很快。在执行过程中,中间人在worker上制作,根据需要复制到其他worker,并尽快丢弃。除非您明确要求,否则永远不会向调度程序或客户端复制副本。

    因此,对于您的问题的根源:您可以通过查看dashboard 来调查性能或您的最佳操作。

    有许多因素决定事情进展的速度:进程可能共享一个 IO 通道;有些任务不会释放 GIL,因此在线程中的并行性很差;组的数量将极大地影响将数据洗牌到组中的数量...另外,调度程序执行的每个任务总是有一些开销。

    由于 Pandas 是高效的,因此对于数据轻松放入内存的情况,它与 Dask 相比表现良好也就不足为奇了。

    【讨论】:

    • 非常感谢。我已经在使用仪表板,但是您在答案中链接的页面中的视频似乎非常有用顺便说一句,我设置中的所有工作人员都在虚拟机上运行并访问一个常见的 nfs 挂载。
    猜你喜欢
    • 2021-09-07
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-02-04
    • 1970-01-01
    • 1970-01-01
    • 2018-07-06
    • 2017-11-23
    相关资源
    最近更新 更多