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
相关产品推荐
相关产品推荐

