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

如何使用PySpark写入HDFS的CSV文件不带.deflate后缀

解决方案

你遇到的问题由两个原因导致:一是压缩参数配置未生效,二是Spark原生CSV写入默认不会主动添加.csv后缀,对应解决方法如下:

1. 修正压缩配置禁用deflate压缩

你传入compression=None未生效,是因为部分Spark版本会优先读取集群默认的输出压缩配置,忽略None值。需要显式在option中指定压缩格式为字符串none,强制覆盖集群默认配置,即可消除.deflate后缀:

from pyspark.sql import SparkSession
spark = SparkSession.builder.getOrCreate()

df = spark.read.parquet("<<path>>")

df.coalesce(1).write.option("header", "false") \
    .option("compression", "none") \
    .mode('overwrite') \
    .csv(output + '/' + end_date)

2. 手动添加.csv后缀

Spark为适配分布式多分片写入场景,默认输出的文件名为part-00000-<随机UUID>格式,不会主动添加.csv后缀,你可以用以下两种方式处理:

  • 方案1:写入完成后重命名(推荐,适合大文件)
    你已经用coalesce(1)输出为单个文件,写入完成后调用HDFS命令重命名文件即可,示例如下:
    import subprocess
    
    output_dir = f"{output}/{end_date}"
    # 查找生成的part文件路径
    part_file_path = subprocess.check_output(
        f"hdfs dfs -ls {output_dir}/part-* | awk '{{print $NF}}'", 
        shell=True
    ).decode().strip()
    # 重命名为带.csv后缀的文件
    subprocess.call(
        f"hdfs dfs -mv {part_file_path} {output_dir}/result.csv", 
        shell=True
    )
    # 可选:删除_SUCCESS、校验文件等无关文件
    subprocess.call(f"hdfs dfs -rm {output_dir}/_SUCCESS {output_dir}/.*.crc", shell=True)
    
  • 方案2:转Pandas写入(仅适合小数据量)
    如果数据量不大,可以将数据拉到Driver端转成Pandas DataFrame后直接写入CSV:
    df.toPandas().to_csv(f"hdfs://{output}/{end_date}/result.csv", header=False, index=False)
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 10:45:04