【发布时间】: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