PySpark JSON写入末尾空行及HDFS存储报错问题求助
PySpark写入大体积JSON文件:解决空行与HDFS报错问题
可行性分析
直接用df.write.json()处理15亿条记录并生成单个JSON文件可行,但存在内存风险:
- 你的代码中使用
collect_list将所有记录聚合到单个Row中,这会把全量数据拉到单个Executor节点,需要该节点有足够内存容纳整个数据集(15亿条记录的结构化数据序列化后体积可能达数十GB甚至上百GB)。 - 若内存不足会直接触发OOM,需提前调整Spark内存参数(如
spark.executor.memory、spark.driver.memory),并建议开启Kryo序列化优化性能。
问题1:消除JSON文件末尾的空行
Spark原生JSON Writer默认会在文件末尾添加换行符,解决方法有两种:
方法1:转成字符串后用Text Writer输出
将整个DataRow序列化为单个JSON字符串,再用write.text()写入,完全控制输出格式:
# 替换原write操作的代码 dfprocessed = dfprocessed.withColumn( "full_json", f.to_json(f.struct("list-item", "version")) ) # 合并为单个分区后写入文本文件 dfprocessed.select("full_json").coalesce(1).write.mode("overwrite").text("./TestJson")
这种方式输出的文件只有一行JSON内容,无额外空行。
方法2:写入后移除空行
如果必须用原生JSON Writer,可在写入后通过文件系统命令删除最后一行:
- 本地文件:
sed -i '$d' ./TestJson/part-*.json - HDFS文件:
hdfs dfs -cat /path/to/part-*.json | sed '$d' | hdfs dfs -put - /path/to/final.json
问题2:HDFS写入报错排查与解决
报错原因需结合具体错误信息分析,常见场景及解决方法:
场景1:权限不足
- 错误信息:
Permission denied: user=xxx, access=WRITE, inode="/somedir_in_HDFS" - 解决:确保执行Spark任务的用户对目标HDFS目录有写入权限,可通过
hdfs dfs -chmod 775 /somedir_in_HDFS调整权限,或切换到有权限的用户执行任务。
场景2:目录已存在且未设置覆盖模式
- 错误信息:
FileAlreadyExistsException: Path /somedir_in_HDFS already exists - 解决:写入时指定
mode("overwrite"),覆盖已有目录:
dfprocessed.write.format("json").mode("overwrite").option("escape", "").save("hdfs:///somedir_in_HDFS")
场景3:内存不足(OOM)
- 错误信息:
OutOfMemoryError或ExecutorLostFailure - 解决:
- 调大Executor内存,提交任务时添加参数:
--executor-memory 64g --driver-memory 32g - 开启Kryo序列化优化:
spark = SparkSession.builder.appName("Test")\ .enableHiveSupport()\ .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer")\ .getOrCreate() - 若聚合后数据量仍过大,需评估是否真的需要单个文件,下游若支持多文件读取,可去掉
coalesce(1),让Spark并行写入多个文件。
- 调大Executor内存,提交任务时添加参数:
场景4:HDFS磁盘空间不足
- 错误信息:
No space left on device - 解决:清理HDFS冗余数据,或扩容磁盘空间。
注意事项
- 15亿条记录聚合到单个分区对内存压力极大,若下游业务允许,优先考虑分文件输出,避免单点内存瓶颈。
- 若必须单个文件,建议在测试环境先使用小数据集验证逻辑,再逐步放大数据量调整内存参数。
内容的提问来源于stack exchange,提问作者Code Heaven
相关产品推荐
相关产品推荐

