PySpark写入单条JSON的报错、空行及Driver内存问题咨询
PySpark单条JSON写入问题分析与解决
一、核心疑问
在PySpark中,若--driver-memory配置值小于待写入的数据量,dataframe.write.save()是否会执行失败?
二、实际场景与问题
我们有一个包含5列、近15亿条记录的DataFrame,需要将其写入为单条JSON文件,目前遇到两个问题:
- 使用
df.write.format('json')写入时,生成的文件为单条记录但存在多余空行; - 使用
df.write.format('json').save('somedir_in_HDFS')写入HDFS时出现报错,Executor容器因外部信号终止(退出码143),Job中止。
示例代码
from pyspark.sql import SparkSession import pyspark.sql.functions as f from pyspark.sql.types import * schema = StructType([ StructField("author", StringType(), False), StructField("title", StringType(), False), StructField("pages", IntegerType(), False), StructField("email", StringType(), False) ]) # 实际为15亿条记录 data = [ ["author1", "title1", 1, "author1@gmail.com"], ["author2", "title2", 2, "author2@gmail.com"], ["author3", "title3", 3, "author3@gmail.com"], ["author4", "title4", 4, "author4@gmail.com"] ] if __name__ == "__main__": spark=SparkSession.builder.appName("Test").enableHiveSupport().getOrCreate() df = spark.createDataFrame(data, schema) dfprocessed=df # 此处包含大量表关联操作 dfprocessed = dfprocessed.agg( f.collect_list(f.struct(f.col('author'), f.col('title'), f.col('pages'), f.col('email'))).alias("list-item")) dfprocessed = dfprocessed.withColumn("version", f.lit("1.0").cast(IntegerType())) dfprocessed.printSchema() dfprocessed.write.format("json").mode("overwrite").option("escape", "").save('./TestJson') # 上述写入会产生多余空行
提交命令
spark-submit --conf spark.port.maxRetries=100 --master yarn --deploy-mode cluster --executor-memory 16g --executor-cores 4 --driver-memory 16g --conf spark.driver.maxResultSize=10g --conf spark.dynamicAllocation.enabled=true --cong spark.dynamicAllocation.maxExecutor=30 --conf spark.dynamicAllocation.minExecutor=15 --code.py
报错信息
"Error occurred in save_daily_report : An error occurred while calling o149.save. :org.apache.spark.sparkException:Job aborted at org.apache.spark.sql.execution.datasource.FileFormatWriter$.write(FileFormatWriter.scala:198) at org.apache.spark.sql.execution.datasource.InsertIntoHadoopFSRelationCommand.run(InsertIntoHadoopFSRelationCommand.scala:159) at org.apache.spark.sql.execution.command.DataWritingCommandExec.sideEffectResult$lzycompute(command.scala:104) ... caused by org.apache.spark.sparkException:Job aborted due to stage Failure:Task 26 in stage 211.0 failed 4 times,most recent failure:Lost task 26.3 in stage 211.0 ExecutorLostFailure(executor xxx exited caused by one the the running task)Reason:Container maked as failed Container exited with a non-zero exit code 143 Killed by external signal Driver stacktrace: at org.apache.spark.scheduler.DAGScheduler.org$apache$spark$scheduler$DAGScheduler$$failedJobAndIndependent Stages (DAGScheduler.scala:1890)
三、失败原因分析
Executor内存不足被YARN杀死:
- 代码中使用
collect_list将15亿条记录聚合为单个列表,导致DataFrame仅1个分区。写入时该分区的Task会在单个Executor上执行,需要加载整个聚合后的数据集到内存,16G的Executor内存无法容纳如此庞大的数据,触发OOM后被YARN NodeManager以外部信号终止(退出码143)。 - 提交命令存在拼写错误:
--cong应为--conf,且spark.dynamicAllocation.maxExecutor应为复数maxExecutors,导致动态资源配置未生效,无法按需扩展Executor资源。
- 代码中使用
多余空行问题:
Spark默认JSON Writer会在每条记录后添加换行符,即使只有一条记录,也会在末尾生成空行,属于组件默认行为。关于driver-memory的疑问解答:
- 普通多分区写入场景:
dataframe.write.save()由Executor执行任务,Driver仅负责提交调度,因此--driver-memory小于数据总量不会导致失败,只要Executor内存足够处理各自分区。 - 单分区聚合写入场景:任务由单个Executor的Task处理,依赖Executor内存而非Driver内存;但若通过
collect()将数据拉到Driver再写入,则需要Driver内存足够容纳整个数据集。
- 普通多分区写入场景:
四、解决方法
1. 解决Executor被杀死(退出码143)问题
- 修正提交命令错误:
将--cong spark.dynamicAllocation.maxExecutor=30改为--conf spark.dynamicAllocation.maxExecutors=30,确保动态资源配置生效。 - 提升Executor内存配置:
增大--executor-memory至32G或更高,同时设置spark.executor.memoryOverhead(如8G),避免堆外内存不足导致容器被杀死。示例:spark-submit --conf spark.port.maxRetries=100 --master yarn --deploy-mode cluster --executor-memory 32g --executor-cores 4 --driver-memory 16g --conf spark.driver.maxResultSize=20g --conf spark.dynamicAllocation.enabled=true --conf spark.dynamicAllocation.maxExecutors=30 --conf spark.dynamicAllocation.minExecutors=15 --conf spark.executor.memoryOverhead=8g --code.py - 优化数据处理逻辑:
- 若下游可接受多行JSON格式,移除
collect_list聚合操作,直接写入多分区文件,避免单Task内存压力。 - 若必须生成单条JSON,可先写入多分区文件,再通过HDFS命令
hdfs dfs -getmerge合并文件,最后手动将多行JSON转换为单个数组格式。
- 若下游可接受多行JSON格式,移除
2. 解决多余空行问题
- 方法一:写入后过滤空行
- 本地文件:使用
sed '/^$/d' input.json > output.json命令去除空行。 - HDFS文件:先下载到本地处理后重新上传,或通过Spark读取文件并过滤空行后写入:
spark.read.text("./TestJson").filter("value != ''").write.text("./TestJson_clean")
- 本地文件:使用
- 方法二:自定义JSON写入逻辑
将聚合后的数据拉到Driver,使用Python原生json模块生成无空行的JSON文件,完全控制输出格式:
注意:此方法需调大Driver内存(如# 聚合后的dfprocessed仅一行 result = dfprocessed.collect()[0].asDict() import json # 写入本地 with open('./TestJson/result.json', 'w') as f: json.dump(result, f) # 写入HDFS可使用pyarrow或hdfs3库 from pyarrow import hdfs fs = hdfs.connect() with fs.open('/path/in/hdfs/result.json', 'w') as f: json.dump(result, f)--driver-memory 32g)及spark.driver.maxResultSize(设为0表示无限制),确保Driver能容纳整个数据集。
内容的提问来源于stack exchange,提问作者Code Heaven
相关产品推荐
相关产品推荐

