如何为读取多文件的Dask DataFrame添加数据源文件名列?
给Dask DataFrame添加数据源文件列
当然可以!这是处理多文件数据时非常实用的需求,Dask提供了几种简单的实现方式,我来逐一介绍:
方法1:使用read_csv的内置参数(推荐)
从Dask 2.16.0版本开始,dd.read_csv()新增了include_path_column参数,能自动为每个分区添加指定列,值就是该分区对应的源文件路径/名称,完全不用手动处理,是最省心的方案:
import dask.dataframe as dd # 读取所有csv文件,并添加名为'partition'的列记录源文件 df = dd.read_csv('file*.csv', include_path_column='partition')
如果只想要文件名而不是完整路径,后续可以用os.path.basename处理这个列:
import os df['partition'] = df['partition'].apply(os.path.basename, meta=('partition', 'str'))
方法2:手动读取并添加列(兼容旧版本或自定义场景)
如果你用的是更早版本的Dask,或者需要更灵活的控制(比如给文件名加前缀、修改格式等),可以用Dask Bag结合自定义读取函数:
import dask.dataframe as dd import pandas as pd from dask.bag import from_sequence import os # 列出所有要读取的文件 file_list = ['file1.csv', 'file2.csv', 'file3.csv'] def read_file_with_source(file_path): # 读取单个文件 df = pd.read_csv(file_path) # 添加源文件列,这里可以自定义处理文件名 df['partition'] = os.path.basename(file_path) return df # 用Bag处理每个文件,再转为Dask DataFrame file_bag = from_sequence(file_list).map(read_file_with_source) df = file_bag.to_dataframe()
方法3:给已有DataFrame的分区添加源信息
如果已经读取了DataFrame,事后想补加源文件列,可以利用Dask分区的属性,结合延迟对象实现:
import dask.dataframe as dd import os from dask.delayed import delayed # 先读取数据 df = dd.read_csv('file*.csv') # 获取每个分区对应的源文件路径 partition_paths = [part.attrs.get('path', '') for part in df.to_delayed()] # 定义给分区添加列的函数 def add_source_column(df_part, source_path): df_part['partition'] = os.path.basename(source_path) return df_part # 给每个分区应用函数,再重新组合成DataFrame delayed_parts = [delayed(add_source_column)(part, path) for part, path in zip(df.to_delayed(), partition_paths)] df_with_source = dd.from_delayed(delayed_parts)
验证结果
最后执行compute()触发计算,就能看到每行都带上了对应的源文件信息:
print(df_with_source.compute())
注意:如果设置了blocksize参数将单个大文件拆分成多个分区,这些子分区的partition列会显示同一个源文件路径,这是符合预期的,因为它们都来自同一个文件。
内容的提问来源于stack exchange,提问作者jpp
相关产品推荐
相关产品推荐

