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

Databricks写入DataFrame生成单个Parquet文件的优化方法咨询

在Databricks中简化CSV转Parquet单文件生成流程

问题场景

在Databricks中用Python将CSV转Parquet时,虽然coalesce(1)能强制生成单分区文件,但Spark会把文件写入单独目录,后续还要手动复制重命名Parquet文件、移动到目标位置并删除原目录,流程繁琐。原代码如下:

mounted_directory_path = "/mnt/myContainer/MyFolder/"
file_list = dbutils.fs.ls(mounted_directory_path)   
    
def insertFirstLineInFile(file_path):
    try:
        header =    ["Model" , "SerialNumber" .....]
        my_df = spark.read.format("csv").option("inferSchema","false").option("sep",",").option("header","false").load(file_path)
        my_df = my_df.toDF(*header)
    
    
        # Extract the filename from file_info.path
        filename = os.path.basename(file_path)
        # rename original filename .csv to .parquet
        if filename.endswith(".csv"): 
            # Replace ".csv" with ".parquet"
            filename = filename[:-4] + ".parquet"
        
        OutputFile = OutputDirectory + filename
        print("Filename:", OutputFile)
        # partition to single file
        my_df_single_partition = my_df.coalesce(1)
        my_df_single_partition.write.option("header", "true").parquet(OutputFile)
        
    except Exception as e:
                print(f"Error reading {file_path}: {str(e)}")

for file_info in file_list:
   insertFirstLineInFile(file_info.path)

优化方案

通过先写入临时目录,再利用dbutils.fs工具自动定位Parquet文件、完成移动重命名并清理临时目录,全程无需手动操作。

优化后代码

mounted_directory_path = "/mnt/myContainer/MyFolder/"
OutputDirectory = "/mnt/myContainer/TargetFolder/"  # 确保目标目录已存在
file_list = dbutils.fs.ls(mounted_directory_path)   

def convert_csv_to_single_parquet(file_path):
    try:
        header = ["Model", "SerialNumber", ...]  # 替换为实际表头字段
        # 读取CSV并指定自定义表头
        my_df = spark.read.format("csv")\
            .option("inferSchema", "false")\
            .option("sep", ",")\
            .option("header", "false")\
            .load(file_path)
        my_df = my_df.toDF(*header)
        
        # 生成目标文件名
        filename = os.path.basename(file_path)
        if filename.endswith(".csv"):
            target_filename = filename[:-4] + ".parquet"
        target_path = f"{OutputDirectory}{target_filename}"
        
        # 创建临时目录存放Spark生成的单分区文件
        temp_dir = f"{OutputDirectory}temp_{target_filename}"
        
        # 将数据写入临时目录
        my_df.coalesce(1).write.option("header", "true").parquet(temp_dir)
        
        # 定位临时目录中的Parquet数据文件(排除_SUCCESS等元文件)
        temp_files = dbutils.fs.ls(temp_dir)
        parquet_file = next((f.path for f in temp_files if f.path.endswith(".parquet")), None)
        
        if parquet_file:
            # 移动并重命名文件到目标路径
            dbutils.fs.mv(parquet_file, target_path)
            # 删除临时目录及其中的冗余文件
            dbutils.fs.rm(temp_dir, recurse=True)
            print(f"成功生成文件: {target_path}")
        else:
            print(f"未找到临时目录中的Parquet文件: {temp_dir}")
            
    except Exception as e:
        print(f"处理文件 {file_path} 时出错: {str(e)}")

# 批量处理目录下的CSV文件
for file_info in file_list:
    if file_info.path.endswith(".csv"):
        convert_csv_to_single_parquet(file_info.path)

关键改动说明

  • 临时目录隔离:先将单分区文件写入临时目录,避免直接写入目标路径产生多余目录结构
  • 自动文件定位:通过dbutils.fs.ls遍历临时目录,精准筛选出实际的Parquet数据文件
  • 自动清理流程:用dbutils.fs.mv完成文件重命名与移动,再用dbutils.fs.rm删除临时目录
  • 增加文件过滤:仅处理CSV格式文件,避免非目标文件触发错误

内容的提问来源于stack exchange,提问作者FunMatters

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 01:52:38