You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

优化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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.22 08:57:40