【问题标题】:parallel execution of dask `DataFrame.set_index()`dask `DataFrame.set_index()` 的并行执行
【发布时间】:2018-11-21 22:44:16
【问题描述】:

我正在尝试在大型 dask 数据帧上创建索引。无论使用哪种调度程序,我都无法使用相当于一个核心的操作。代码是:

(ddf.
 .read_parquet(pq_in)
 .set_index('title', drop=True, npartitions='auto', shuffle='disk', compute=False)
 .to_parquet(pq_out, engine='fastparquet', object_encoding='json', write_index=True, compute=False)
 .compute(scheduler=my_scheduler)
)

我在一台 64 核机器上运行它。我能做些什么来利用更多的核心?还是set_index 本身就是顺序的?

【问题讨论】:

    标签: dataframe concurrency parallel-processing dask dask-distributed


    【解决方案1】:

    这应该使用多个内核,尽管使用磁盘进行洗牌可能会引入其他瓶颈,例如本地硬盘驱动器。您通常不受额外 CPU 内核的约束。

    在您的情况下,我会在单台机器上使用分布式调度程序,以便您可以使用诊断仪表板更深入地了解您的计算。

    【讨论】:

    • 使用分布式调度程序进行更改并设置shuffle='disk' 提高了并行性,但似乎使 dask 尝试将所有数据加载到内存中。是否可以对大于内存的数据进行并行洗牌?
    • 实际上我的数据确实适合内存。问题是分布式调度程序似乎将整个数据集加载到每个工作进程中。
    猜你喜欢
    • 1970-01-01
    • 2020-10-02
    • 2018-11-14
    • 2017-05-13
    • 1970-01-01
    • 1970-01-01
    • 2018-03-19
    • 2019-04-19
    • 2022-11-10
    相关资源
    最近更新 更多