【问题标题】:Python and Dask - reading and concatenating multiple filesPython 和 Dask - 读取和连接多个文件
【发布时间】:2020-08-07 03:08:22
【问题描述】:

我有一些parquet 文件,它们都来自同一个域,但结构有所不同。我需要连接所有这些。下面是这些文件的一些示例:

file 1:
A,B
True,False
False,False

file 2:
A,C
True,False
False,True
True,True

我要做的是以最快的方式读取和连接这些文件,获得以下结果:

A,B,C
True,False,NaN
False,False,NaN
True,NaN,False
False,NaN,True
True,NaN,True

为此,我使用以下代码,使用 (Reading multiple files with Dask, Dask dataframes: reading multiple files & storing filename in column) 提取:

import glob

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

def read_parquet(path):
    return pd.read_parquet(path)

if __name__=='__main__':

    files = glob.glob('test/*/file.parquet')

    print('Start dask client...')
    client = Client()

    results = [dd.from_delayed(dask.delayed(read_parquet)(diag)) for diag in diag_files]

    results = dd.concat(results).compute()

    client.close()

这段代码有效,它已经是我能想到的最快的版本(我尝试了顺序pandas 和multiprocessing.Pool)。我的想法是 Dask 可以理想地开始连接的一部分,同时仍然读取一些文件,但是,从任务图中,我看到每个 parquet 文件的元数据的一些顺序读取,请参见下面的屏幕截图:

任务图的第一部分是read_parquet 和read_metadata 的混合。第一部分始终只显示执行的 1 个任务(在任务处理选项卡中)。第二部分是from_delayed 和concat 的组合,它使用了我所有的工人。

关于如何加快文件读取并减少图表第一部分的执行时间有什么建议吗?

【问题讨论】:

    标签: python dask dask-delayed


    【解决方案1】:

    您的代码的问题是您使用 Pandas 版本的 read_parquet.

    改为使用:

    • dask 版本的 read_parquet,
    • 客户端提供的map和gather方法,
    • dask 版本的 concat,

    类似:

    def read_parquet(path):
        return dd.read_parquet(path)
    
    def myRead():
        L = client.map(read_parquet, glob.glob('file_*.parquet'))
        lst = client.gather(L)
        return dd.concat(lst)
    
    result = myRead().compute()
    

    在此之前我创建了一个客户端,只有一次。 原因是在我早期的实验中,我遇到了一个错误 当我尝试再次创建它(在函数中)时的消息,甚至 虽然第一个实例之前已经关闭。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-04-14
      • 1970-01-01
      • 1970-01-01
      • 2014-04-12
      • 1970-01-01
      相关资源
      最近更新 更多