【发布时间】:2020-03-16 13:30:31
【问题描述】:
我有大量 csv 文件(file_1.csv、file_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.csv、file_2.csv)
仪器10 的最终文件应如下所示,其中t7、t9 的观察结果与其他文件中仪器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会连接列值 -
value 是
str是否正常,或者您最终需要转换为浮动? 时间和仪器呢?转换类型将帮助您节省大量内存。