【问题标题】:Slow Dask performance on CSV date parsing?CSV 日期解析的 Dask 性能缓慢?
【发布时间】:2017-01-15 15:02:52
【问题描述】:

我一直在对一大堆文件进行大量文本处理,包括大型 CSV 文件和大量小型 XML 文件。有时我在做汇总计数,但很多时候我在做 NLP 类型的工作,以便更深入地了解这些文件中的内容,而不是标记或已经结构化的内容。

我一直在使用多处理库在多个 CPU 上执行这些计算,但我爱上了 Dask 背后的想法,它在网上和同事都得到了强烈推荐。

我在这里问了一个关于 Dask 性能的类似问题:

Slow Performance with Python Dask bag?

和 MRocklin (https://stackoverflow.com/users/616616/mrocklin) 告诉我,加载大量小文件可能会破坏性能。

然而,当我在单个大文件 (200mb) 上运行它时,我仍然不能让它表现得很好。这是一个例子:

我有一个 900,000 行的 CSV 推文文件,我想快速加载它并解析“created_at”字段。以下是我完成的三种方法以及每种方法的基准。我在具有 16GB 内存的新 i7 2016 MacBook Pro 上运行此程序。

import pandas
import dask.dataframe as dd
import multiprocessing

%%time
# Single Threaded, no chunking
d = pandas.read_csv("/Users/michaelshea/Documents/Data/tweet_text.csv", parse_dates = ["created_at"])
print(len(d))

CPU时间:用户2分31秒,系统:807毫秒,总计:2分32秒 挂壁时间:2分32秒

%%time
# Multithreaded chunking
def parse_frame_dates(frame):
    frame["created_at"] = pandas.to_datetime(frame["created_at"])
    return(frame)

d = pandas.read_csv("/Users/michaelshea/Documents/Data/tweet_text.csv", chunksize = 100000)
frames = multiprocessing.Pool().imap_unordered(get_count, d)
td = pandas.concat(frames)
print(len(td))

CPU 时间:用户 5.65 秒,系统:1.47 秒,总计:7.12 秒 挂墙时间:1分10秒

%%time
# Dask Load
d = dd.read_csv("/Users/michaelshea/Documents/Data/tweet_text.csv", 
                 parse_dates = ["created_at"], blocksize = 10000000).compute()

CPU时间:用户2分59秒,系统:26.2秒,总计:3分25秒 挂墙时间:3分12秒

我在许多不同的 Dask 比较中发现了这些类型的结果,但即使让它正常工作也可能为我指明正确的方向。

简而言之,我怎样才能让 Dask 在这些任务中发挥最佳性能?为什么它的性能似乎不如其他方式的单线程和多线程技术?

【问题讨论】:

    标签: python multithreading performance pandas dask


    【解决方案1】:

    我怀疑 Pandas read_csv 日期时间解析代码是纯 python,因此不会从使用线程中受益,这是 dask.dataframe 默认使用的。

    您可能会在使用进程时看到更好的性能。

    我怀疑以下方法会更快:

    import dask.multiprocessing
    dask.set_options(get=dask.multiprocessing.get)  # set processes as default
    
    d = dd.read_csv("/Users/michaelshea/Documents/Data/tweet_text.csv", 
                    parse_dates = ["created_at"], blocksize = 10000000)
    len(d)
    

    进程的问题是进程间通信会变得昂贵。我在上面明确计算len(d) 而不是d.compute(),以避免不得不拾取工作进程中的所有pandas 数据帧并将它们移动到主调用进程。在实践中这很常见,因为人们很少需要完整的数据帧,而是对数据帧进行一些计算。

    这里的相关文档是http://dask.readthedocs.io/en/latest/scheduler-choice.html

    您可能还想在单台机器上使用distributed scheduler,而不是使用多处理调度程序。上面引用的文档中也对此进行了描述。

    $ pip install dask distributed
    
    from dask.distributed import Client
    c = Client()  # create processes and set as default
    
    d = dd.read_csv("/Users/michaelshea/Documents/Data/tweet_text.csv", 
                    parse_dates = ["created_at"], blocksize = 10000000)
    len(d)
    

    【讨论】:

    • 计时方法,OP不一样。传递parse_dates=... 是一种相当稳健的方法,但我不得不退回到较慢的解析(在python中)。您几乎总是希望使用 .to_datetime 简单地读取 csv、THEN、后处理,特别是您可能需要使用 format= 参数或其他选项,具体取决于日期。 ,YMMV。事实上,这种方法对 dask 非常友好,仅供参考(因为它们是分开的,虽然是连续的任务)。
    • 谢谢,马修!这产生了巨大的变化。在我的示例中,转到多进程时,将其从单线程的 2 分钟降至 Dask 的 30 秒。更符合我的预期。我将深入研究有关调度程序选择和分布式调度程序使用的信息。我显然没有做完所有的功课。再次感谢!
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2017-08-20
    • 2019-10-18
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-02-07
    • 2015-09-01
    相关资源
    最近更新 更多