【问题标题】:Why does running compute() on a filtered Dask dataframe take so long?为什么在过滤后的 Dask 数据帧上运行 compute() 需要这么长时间?
【发布时间】:2020-03-16 20:10:27
【问题描述】:

我正在使用以下方式读取数据: ddf1 = dd.read_sql_table('mytable', conn_string, index_col='id', npartitions=8)

当然,由于惰性计算,这会立即运行。这个表有几亿行。

接下来,我要过滤这个 Dask 数据框:

ddf2 = ddf1.query('some_col == "converted"')

最后,我想将其转换为 Pandas 数据框。结果应该只有大约 8000 行:

ddf3 = ddf2.compute()

但是,这需要很长时间(约 1 小时)。我可以就如何大幅加快速度获得任何建议吗?我试过使用.compute(scheduler='threads'),改变分区的数量,但到目前为止没有一个工作。我做错了什么?

【问题讨论】:

  • 大概是因为它做了很多工作吧?
  • 如果我有错误的印象,请原谅我,但我认为 Dask 应该大幅加快速度?
  • 可能。你把它比作什么替代品?几亿行是很多数据。你的设置到底是什么?计算集群?还是你的笔记本电脑? dask 不是魔法
  • 我将其与通过 Pandas 将整个表加载到内存中进行比较。我的设置是我的笔记本电脑,它有四个内核,但我也在具有高内存的 EC2 实例上尝试过这个,我仍然表现出类似的性能问题,这让我相信我没有正确进行配置
  • 好吧,只用 pandas 需要多长时间?您能否为您的问题添加更多详细信息?

标签: python pandas parallel-processing dask dask-dataframe


【解决方案1】:

首先,您可以使用 sqlalchemy 表达式语法对查询中的过滤子句进行编码,并在服务器端进行过滤。如果数据传输是您的瓶颈,那是您最好的解决方案,尤其是过滤列被索引。

根据您的数据库后端,sqlalchemy 可能不会释放 GIL,因此您的分区不能在线程中并行运行。你得到的只是线程之间的争用和额外的开销。您应该将distributed scheduler 与进程一起使用。

当然,看看你的CPU和内存使用情况;使用分布式调度程序,您还可以访问诊断仪表板。您还应该关注每个分区在内存中的大小。

【讨论】:

  • 感谢您的回复。我不能过滤服务器端的原因是因为我有数千个过滤操作要执行,每个过滤操作都在巨大表的不同子集上。我应该明确表示我不仅只进行一次过滤。这个想法是将整个表加载到 Dask 数据帧中,并根据需要对 Dask 数据帧执行单独的过滤器,而不是查询数据库数千次并给服务器带来压力。实际上,我确实尝试过使用远程 EC2 服务器上的进程的分布式调度程序,它表现出类似的糟糕性能。
  • 您对划分核心数、线程数、worker数等有什么建议吗?
猜你喜欢
  • 1970-01-01
  • 2015-04-15
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-09-22
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多