如何高效为Dask DataFrame添加文件名列以分组绘制时序图?
用Dask高效为每个CSV/Parquet文件添加文件名列并分组
问题背景
我有大约400个包含多变量时间序列的CSV文件(每个文件含时间列及其他多变量列),最终目标是选取部分变量并绘制400条时序图。用Dask读取文件后发现,必须按数据来源的文件名分组才能绘制每条文件对应的时序图,而非合并后的单条。想找高效的Dask方法为每个文件添加文件名列,Parquet文件也可作为备选。
之前尝试的代码存在问题:
import dask.dataframe as dd import os filenames = ['part0.parquet', 'part1.parquet', 'part2.parquet'] df = dd.read_parquet(filenames, engine='pyarrow') df = df.assign(file=lambda x: filenames[x.index]) df_grouped = df.groupby('file')
已知from_delayed方法可行,但会丢失并行计算能力,希望得到更优方案。
解决方案
方法1:用Dask内置参数自动添加文件路径列
Dask的read_parquet和read_csv都自带include_path_column参数,能直接为每个文件添加包含文件路径的列,全程保留并行计算能力,是最优方案。
处理Parquet文件
import dask.dataframe as dd filenames = ['part0.parquet', 'part1.parquet', 'part2.parquet'] # 读取时指定include_path_column,自动生成file_path列 df = dd.read_parquet(filenames, engine='pyarrow', include_path_column='file_path') # 若只需文件名而非完整路径,用Dask字符串方法提取 df['file'] = df['file_path'].str.split('/').str[-1] # 按文件名分组 df_grouped = df.groupby('file')
处理CSV文件
import dask.dataframe as dd csv_filenames = ['data1.csv', 'data2.csv', ...] # 读取CSV时同样启用include_path_column df = dd.read_csv(csv_filenames, include_path_column='file_path') # 提取文件名 df['file'] = df['file_path'].str.split('/').str[-1] df_grouped = df.groupby('file')
方法2:手动为分区绑定文件名(兼容旧版Dask)
如果你的Dask版本不支持include_path_column,可以通过map_partitions结合分区元数据添加文件名,同样不损失并行性:
import dask.dataframe as dd import os filenames = ['part0.parquet', 'part1.parquet', 'part2.parquet'] df = dd.read_parquet(filenames, engine='pyarrow') # 定义给分区添加文件名的函数 def add_filename(partition, filename): partition['file'] = filename return partition # 遍历每个分区,绑定对应的文件名后合并 df = dd.multi.concat( [df.get_partition(i).map_partitions(add_filename, filename=filenames[i]) for i in range(df.npartitions)] ) df_grouped = df.groupby('file')
原代码错误原因
原代码里lambda x: filenames[x.index]逻辑完全错误:x.index是数据的行索引,和文件在filenames列表中的索引没有任何对应关系,这会导致把行索引值当成列表下标去取文件名,结果完全不符合预期。
内容的提问来源于stack exchange,提问作者Ben
相关产品推荐
相关产品推荐

