ADLS中读取Parquet生成单CSV并删除_delta_log文件夹
解决方案
问题原因
Spark的write.save()方法默认会将输出写入一个目录,同时生成SUCCESS、_COMMITTED等元数据文件,这是分布式计算框架的特性,无法直接输出单个指定文件名的CSV到目标路径。另外你还需要过滤掉_delta_log文件夹的内容,并最终删除该文件夹。
分步实现代码(Python)
1. 读取Parquet文件(排除_delta_log)
通过pathGlobFilter参数只读取.parquet后缀的文件,自动排除_delta_log文件夹:
target_path = "/mnt/adls/2022-12-05" # 读取目标路径下所有Parquet文件,跳过_delta_log df = spark.read.parquet(f"{target_path}/*", pathGlobFilter="*.parquet")
2. 生成指定文件名的单个CSV
先将数据写入临时目录,再提取其中的分片CSV文件,重命名为csv_file.csv后移动到目标路径,最后清理临时文件:
import os import shutil # 临时输出目录,用于存放Spark生成的中间文件 temp_output_dir = f"{target_path}/temp_csv" # 写入临时目录,coalesce(1)确保仅生成一个分片CSV df.coalesce(1).write.format("csv")\ .mode("overwrite")\ .option("header", "true")\ # 根据需求决定是否保留表头 .save(temp_output_dir) # 查找临时目录中的CSV文件 csv_part_files = [f for f in os.listdir(temp_output_dir) if f.endswith(".csv")] if csv_part_files: # 将分片CSV重命名并移动到目标路径 source_csv = os.path.join(temp_output_dir, csv_part_files[0]) target_csv = os.path.join(target_path, "csv_file.csv") shutil.move(source_csv, target_csv) # 删除临时目录 shutil.rmtree(temp_output_dir)
3. 删除_delta_log文件夹
处理完CSV后,直接删除目标路径下的_delta_log文件夹:
delta_log_dir = os.path.join(target_path, "_delta_log") if os.path.exists(delta_log_dir): shutil.rmtree(delta_log_dir)
注意事项
- 如果数据量极大,
coalesce(1)会将所有数据集中到单个Executor,可能引发性能问题或内存溢出,此时需评估是否真的需要单个CSV文件。 - 确保执行代码的账号拥有ADLS目标路径的读写、创建/删除目录的权限。
- 若使用Scala实现,文件操作可借助
java.nio.file包的API,逻辑与Python一致。
内容的提问来源于stack exchange,提问作者Nilesh Patel
相关产品推荐
相关产品推荐

