【问题标题】:How to separate files using dask groupby on a column如何在列上使用 dask groupby 分隔文件
【发布时间】:2020-03-16 13:30:31
【问题描述】:

我有大量 csv 文件(file_1.csvfile_2.csv),按时间段分隔,无法放入内存。每个文件将采用下面提到的格式。


| instrument | time | code     | val           |
|------------|------|----------|---------------|
| 10         | t1   | c1_at_t1 | v_of_c1_at_t1 |
| 10         | t1   | c2_at_t1 | v_of_c2_at_t1 |
| 10         | t2   | c1_at_t2 | v_of_c1_at_t2 |
| 10         | t2   | c3_at_t2 | v_of_c3_at_t2 |
| 11         | t1   | c4_at_t1 | v_of_c4_at_t1 |
| 11         | t1   | c5_at_t1 | v_of_c5_at_t1 |
| 12         | t2   | c6_at_t2 | v_of_c6_at_t2 |
| 13         | t3   | c9_at_t3 | v_of_c9_at_t3 |

每个文件都与格式一致的仪器日志有关。有一组仪器可以在给定的时间戳(time)发出不同的代码(code)。给定仪器的给定time 处的code 值保存在val 列中

我想使用instrument 列(例如:10)拆分每个文件(例如:file_1.csv),然后在所有文件中加入为仪器提取的文件(例如:10)( file_1.csv, file_2.csv)

我正在考虑在instrument 列上使用dask groupby 操作。是否有任何替代或更好的方法来代替使用groupby 或通过instrument 提取文件的更好方法?

我为执行上述操作而编写的代码是

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

client = Client()

def read_files(files):

    files = glob.glob(files)

    for f in files:

        df = dd.read_csv(f, blocksize='256MB')
        unique_inst = df['instrument'].unique()
        gb = df.groupby('instrument')  

        for v in unique_inst:
            gb.get_group(v).to_parquet(f'{v}_{f[:-4]}.parquet')

    pass

获得f'{v}_{f[:-4]}.parquet' 格式的文件后,我可以使用从所有文件中提取的pandas 连接它们(file_1.csvfile_2.csv

仪器10 的最终文件应如下所示,其中t7t9 的观察结果与其他文件中仪器10 的观察结果相连接

time | code     | val           |
-----|----------|---------------|
t1   | c1_at_t1 | v_of_c1_at_t1 |
t1   | c2_at_t1 | v_of_c2_at_t1 |
t2   | c1_at_t2 | v_of_c1_at_t2 |
t2   | c3_at_t2 | v_of_c3_at_t2 |
t7   | c4_at_t7 | v_of_c4_at_t7 |
t9   | c5_at_t9 | v_of_c5_at_t9 |

【问题讨论】:

  • 我不清楚像file1.csv 这样的单个文件是否适合内存?
  • 没有一个csv 文件可以放入内存,因为每个文件的大小超过 100 GB
  • 哇,这有点改变。问题: 1. 数据类型是什么? 2. 只有一个文件,其中df 是一个dask.dataframe 是df.groupby(""instrument")["val"].sum().compute() 给你一个内存错误?
  • 对不起,如果我在问题中的第一个陈述引起任何混淆,因为我提到了大型 csv 文件 cant be fit into memory。所有列值都作为字符串存储在 csvs 中。我可以做df.groupby(""instrument")["val"].sum().compute(),但是每一列都是一个字符串,sum 会连接列值
  • valuestr 是否正常,或者您最终需要转换为浮动? 时间仪器呢?转换类型将帮助您节省大量内存。

标签: python pandas dask


【解决方案1】:

如果每个文件都适合内存,你可以试试这个:

import dask.dataframe as dd
import pandas as pd
import numpy as np
import os

生成虚拟文件

fldr_in = "test_in"
fldr_out = "test_out"

N = int(1e6)
for i in range(10):
    fn = f"{fldr_in}/file{i}.csv"
    os.makedirs(os.path.dirname(fn), exist_ok=True)
    df = pd.DataFrame({"instrument":np.random.randint(10,100,N),
                       "value":np.random.rand(N)})
    df.to_csv(fn, index=False)

定义函数

以下函数为路径fldr_out/instrument=i/fileN.csv中的每个乐器保存到镶木地板

def fun(x, fn, fldr_out):
    inst = x.instrument.unique()[0]
    filename = os.path.basename(fn)
    fn_out = f"{fldr_out}/instrument={inst}/{filename}"
    fn_out = fn_out.replace(".csv", ".parquet")
    os.makedirs(os.path.dirname(fn_out), exist_ok=True)
    x.drop("instrument", axis=1)\
     .to_parquet(fn_out, index=False)

您可以将它与 group by 一起使用

for f in files:
    fn = f"{fldr_in}/{f}"
    df = pd.read_csv(fn)
    df.groupby("instrument").apply(lambda x: fun(x, fn, fldr_out))

使用 dask 进行分析

现在您可以使用dask 来读取结果并执行您的分析

df = dd.read_parquet(fldr_out)

【讨论】:

    【解决方案2】:

    我不确定您需要达到什么目标,但我认为您不需要任何 group by 来解决您的问题。在我看来,这是一个简单的过滤问题。

    您可以遍历所有文件并创建新的仪器文件并附加到这些文件上。

    我也没有示例文件可供试验,但我认为您也可以使用带有 chunksize 的 pandas 来读取大型 csv 文件。

    例子:

    import pandas as pd
    import glob
    import os
    
    # maybe play around to get better performance 
    chunksize = 1000000
    
    files = glob.glob('./file_*.csv')
    for f in files:
    
         for chunk in pd.read_csv(f, chunksize=chunksize):
             u_inst = chunk['instrument'].unique()
    
             for inst in u_inst:
                 # filter instrument data
                inst_df = chunk[chunk.instrument == inst]
                # filter columns
                inst_df = inst_df[['time', 'code', 'val']]
                # append to instrument file
                # only write header if not exist yet
                inst_file = f'./instrument_{inst}.csv'
                file_exist = os.path.isfile(inst_file)
                inst_df.to_csv(inst_file, mode='a', header=not file_exist)
    

    【讨论】:

    • 谢谢。这种方法的一个警告是,我们正在遍历每个块,但是u_inst 的数量是可用的。您对阅读每一行并将其附加到正确的 inst_file 有何看法?或使用上面提到的与chunksize1 相同的方法
    • 此类操作的瓶颈通常是IO。更大的块大小允许一次读取更多内容到您的内存中。块上的过滤本身是一种查找,并且速度非常快。因此,我建议使用大块大小。也许我提出的更大。如果 chunksize 1 会产生好的结果,我会感到惊讶。
    • 我用 chunksize =1 尝试了同样的方法。与分块读取然后通过inst过滤相比,该过程非常慢
    猜你喜欢
    • 1970-01-01
    • 2021-04-29
    • 1970-01-01
    • 1970-01-01
    • 2016-03-19
    • 2017-11-30
    • 1970-01-01
    • 2018-01-08
    • 2018-06-18
    相关资源
    最近更新 更多