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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 11:01:36