【问题标题】:using Dask to load many CSV files with different columns使用 Dask 加载具有不同列的许多 CSV 文件
【发布时间】:2021-10-09 08:40:05
【问题描述】:

我在 AWS s3 中保存了许多 CSV 文件,其中包含相同的第一组列和许多可选列。我不想一个一个地下载它们,而是使用pd.concat 来阅读它,因为这需要很多时间并且必须适合计算机内存。相反,我正在尝试使用 Dask 来加载和汇总所有这些文件,而应将可选列视为零。

如果我可以使用相同的所有列:

    import dask.dataframe as dd
    addr = "s3://SOME_BASE_ADDRESS/*.csv"
    df = dd.read_csv(addr)
    df.groupby(["index"]).sum().compute()

但它不适用于列数不同的文件,因为 Dask 假设它可以将第一列用于所有文件:

文件“.../lib/python3.7/site-packages/pandas/core/internals/managers.py”,第 155 行,在 set_axis '值有 {new} 个元素'.format(old=old_len, new=new_len)) ValueError:长度不匹配:预期轴有 64 个元素,新值有 62 个元素

根据this thread,我可以预先读取所有标题(例如,在我生成并保存所有小 CSV 时写入它们)或使用类似这样的东西:

df = dd.concat([dd.read_csv(f) for f in filelist])

我想知道这个解决方案是否真的比直接使用 pandas 更快/更好?一般来说,我想知道解决此问题的最佳(主要是最快)方法是什么?

【问题讨论】:

    标签: pandas csv amazon-s3 dask dask-distributed


    【解决方案1】:

    在将数据帧转换为 dask 数据帧之前,最好使用 delayed 对其进行标准化(很难判断这是否最适合您的用例)。

    import dask.dataframe as dd
    from dask import delayed
    
    list_files = [...] # create a list of files inside s3 bucket
    list_cols_to_keep = ['col1', 'col2']
    
    @delayed
    def standard_csv(file_path):
        df = pd.read_csv(file_path)
        df = df[list_cols_to_keep]
        # add any other standardization routines, e.g. dtype conversion
        return df
    
    ddf = dd.from_delayed([standard_csv(f) for f in list_files])
    

    【讨论】:

    • 如果我理解正确dask.dataframe 应该已经使用delayed 选项。意思是您提供的内容与我提供的内容有些相同:df = dd.concat([dd.read_csv(f) for f in filelist]) 当两者都应跟随df.groupby(["index"]).sum().compute() 以获得我的最终结果。我说的对吗?
    • 我知道这可能不太清楚:我想要最终 df 中的所有列
    • 在示例中,list_cols_to_keep 可能是所有可用列,按照您想要的顺序。
    • 但我事先不知道可能的选项列表(太多了,我只想要现有的)
    【解决方案2】:

    我最终放弃了使用Dask,因为它太慢了,使用aws s3 sync 下载数据并使用multiprocessing.Pool 读取和连接它们:

    # download:
    def sync_outputs(out_path):
        local_dir_path = f"/tmp/outputs/"
        safe_mkdir(os.path.dirname(local_dir_path))
        cmd = f'aws s3 sync {job_output_dir} {local_dir_path} > /tmp/null' # the last part is to avoid prints
        os.system(cmd)
        return local_dir_path
    
    # concat:
    def read_csv(path):
        return pd.read_csv(path,index_col=0)
    
    def read_csvs_parallel(local_paths):
        from multiprocessing import Pool
        import os
        with Pool(os.cpu_count()) as p:
            csvs = list(tqdm(p.imap(read_csv, local_paths), desc='reading csvs', total=len(paths)))
        return csvs
    
    # all togeter:
    def concat_csvs_parallel(out_path):
        local_paths = sync_outputs(out_path)
        csvs = read_csvs_parallel(local_paths)
        df = pd.concat(csvs)
        return df
    

    aws s3 sync 在大约 30 秒内下载了大约 1000 个文件(每个约 1KB),并且读取比使用多进程(8 核)需要 3 秒,这比使用multiprocessing 下载文件要快得多(几乎 2 分钟1000 个文件)

    【讨论】:

      猜你喜欢
      • 2019-04-09
      • 2019-10-11
      • 2014-10-12
      • 2022-12-12
      • 2020-07-16
      • 2021-10-04
      • 1970-01-01
      • 2013-09-25
      • 2023-01-18
      相关资源
      最近更新 更多