【问题标题】:dask loading multiple parquet files with different column selectionsdask 加载具有不同列选择的多个镶木地板文件
【发布时间】:2019-10-11 09:46:52
【问题描述】:

我想使用 Dask 从存储在不同目录中的许多 parquet 文件中加载特定列,并且每个 parquet 需要加载不同的列。我想使用 Dask,以便我可以在单台机器上使用多个内核。我知道如何将文件列表或通配符传递给dd.read_parquet 以指示多个文件(例如*.parquet),但我看不到为每个文件传递要读取的不同列集的方法。我想知道这是否可以使用dask.delayed 来完成。

我的具体情况是:

我将大型单细胞基因表达数据集(约 30,000 行/基因乘以约 10,000 列/细胞)作为 parquet 文件存储在不同的目录中。每个目录有两个 parquet 文件 1) 大型基因表达数据(单元格作为列)和 2) 单元格元数据(单元格作为行,单元格元数据作为列)。我正在使用较小的元数据拼花文件来查找较大文件中需要的列/单元格。例如,我将使用元数据 parquet 文件查找特定单元格类型的所有单元格,然后仅从较大的文件中加载这些单元格。我可以使用 Pandas 做到这一点,但我想使用 Dask 进行并行处理。

【问题讨论】:

    标签: pandas dask parquet dask-distributed


    【解决方案1】:

    如果您可以使用Pandas .read_parquet 和specifying columns 来执行此操作(请参阅最后一个代码示例),那么一种可能的方法是通过替换来延迟您现有的 Pandas 特定方法

    pd.read_parquet(..., columns=[list_of_cols])
    

    通过

    dask.delayed(pd.read_parquet)(..., columns=[list_of_cols])
    

    正如你所建议的那样。

    编辑

    我必须对成对的 .csv 文件的单个目录执行类似的操作 - 元数据和相应的光谱。我的过滤逻辑是最小的,所以我创建了一个 Python 字典,其键是元数据逻辑(生成文件名),值是列列表。我遍历了字典键值对巴黎和

    • 使用dd.read_csv(..., columns=[list_of_cols])从相关光谱文件中读取相应的列列表
    • 将ddf 附加到一个空白列表(显然后面跟着dd.concat() 在循环之后将它们垂直连接在一起)

    不过,就我而言,元数据内容以可预测的方式发生变化,这就是为什么我可以使用dictionary comprehension 以编程方式组装字典的原因。

    【讨论】:

    • 如果 read_parquet 不适合内存怎么办?
    • 在加载文件时尝试过滤数据。 pd.read_parquet() 接受 filters 参数。也许这可以用来减少数据的大小,使其适合内存?查看示例here。
    • 去过那里:它在某些地方有帮助,但如果生成的数据框也太大而无法放入内存,那么我仍然会对 dask 感兴趣。
    • 我想我引起了一些混乱 - 在我之前的评论中,我假设 pd.read_parquet 正在被使用。如果您的数据适合内存,那么pd.read_parquet 可以正常工作。为混乱道歉。
    • 感谢@edesz。我认为这是要走的路!只是一个警告:目前正在处理这样的事情,但如果您的源由几千个镶木地板文件组成,这并不容易。清理和重新分区确实是这里有趣的部分。我们目前的工作流程为此大约生成了 100 万个任务,这真的很烦人,真的很困扰调度程序。 load-clean-analysis 图只是通过使用构建 read_parquet 来变大。有趣的是,这种方法也可能达到其极限。
    猜你喜欢
    • 2020-01-06
    • 2020-09-09
    • 2022-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-01-29
    • 2021-07-05
    • 2023-03-19
    相关资源
    最近更新 更多