如何使用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
相关产品推荐
相关产品推荐

