优化Dask分区首行计算,实现Parquet按原CSV文件名命名
问题描述
我的目标是读取多个CSV文件,完成计算后用to_parquet的partition_on选项保存为Parquet数据集。受内存限制,保存前无法重新索引和重分区,要求每个CSV文件对应一个独立Parquet分区文件,且不能使用默认文件名(如part.0.parquet),避免后续新增文件时出现重名问题,因此希望用原CSV文件名命名对应的Parquet文件。
当前实现方式:读取CSV时添加orig_file_name列(同一分区内所有行的文件名相同),通过map_partitions获取各分区首行的文件名,调用compute()得到文件名列表,再使用name_function为Parquet文件命名。该方案可实现需求,但compute()操作耗时过长,请教如何限制计算仅针对各分区首行以提升效率?
当前实现代码
def get_first_element(partition): return partition['orig_file_name'].iloc[0] first_elements = ddf.map_partitions(get_first_element).compute() def name_function(part_idx): return f"{first_elements[part_idx]}.parquet" ddf.to_parquet(path=target_directory, engine='pyarrow', partition_on=['date', 'hour'], name_function=name_function, write_index=True)
补充复现代码
@dask.delayed def process(file_path): df = pd.DataFrame({'col1':[0, 1, 2, 3], 'col2':[4, 5, 6, 7], 'col3':[88, 88, 99, 99]}) # 实际代码中是read_csv file_name = 'aaa' df.to_parquet(f'{file_name}.parquet', partition_cols=['col3']) dask.compute(*[process(f) for f in [1]])
解决方案
优化核心:避免触发全量计算,仅提取必要的文件名信息
方法1:用persist()替代compute(),仅加载文件名数据
修改分区处理函数,返回仅包含首行文件名的小型数据集,通过persist()先缓存这部分数据,再转换为列表,避免全量计算整个DataFrame:
def get_first_element(partition): # 返回单元素Series,保留分区结构 return pd.Series([partition['orig_file_name'].iloc[0]]) # persist()仅加载文件名相关数据,不触发全量计算 first_elements = ddf.map_partitions(get_first_element).persist() # 此时compute()仅计算文件名部分,耗时大幅降低 first_elements_list = first_elements.compute() def name_function(part_idx): return f"{first_elements_list[part_idx]}.parquet" ddf.to_parquet(path=target_directory, engine='pyarrow', partition_on=['date', 'hour'], name_function=name_function, write_index=True)
方法2:读取阶段直接建立分区与文件名的映射(最优方案)
如果每个CSV文件对应一个Dask分区(即未设置blocksize拆分文件),可以直接在读取时记录文件列表,跳过从DataFrame中提取文件名的步骤:
from dask.dataframe import read_csv # 提前准备所有CSV文件路径 file_paths = ["path/to/file1.csv", "path/to/file2.csv", ...] # blocksize=None确保每个文件对应一个分区 ddf = read_csv(file_paths, blocksize=None) # 生成文件名列表(去掉路径和.csv后缀) file_names = [path.split('/')[-1].replace('.csv', '') for path in file_paths] def name_function(part_idx): return f"{file_names[part_idx]}.parquet" ddf.to_parquet(path=target_directory, engine='pyarrow', partition_on=['date', 'hour'], name_function=name_function, write_index=True)
这种方法完全避免了从DataFrame中提取数据的操作,效率最高。
注意事项
- 若CSV文件被
blocksize拆分为多个分区,方法2不适用,优先选择方法1。 - 使用
partition_on时,单个原始CSV分区可能会生成多个子分区(对应不同的date/hour值),此时name_function的part_idx是子分区索引,可在文件名后添加子分区后缀(如file1_0.parquet)避免重名。
内容的提问来源于stack exchange,提问作者barbarab
相关产品推荐
相关产品推荐

