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

PySpark写入单条JSON的报错、空行及Driver内存问题咨询

PySpark单条JSON写入问题分析与解决

一、核心疑问

在PySpark中,若--driver-memory配置值小于待写入的数据量,dataframe.write.save()是否会执行失败?

二、实际场景与问题

我们有一个包含5列、近15亿条记录的DataFrame,需要将其写入为单条JSON文件,目前遇到两个问题:

  1. 使用df.write.format('json')写入时,生成的文件为单条记录但存在多余空行;
  2. 使用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)

三、失败原因分析

  1. Executor内存不足被YARN杀死:

    • 代码中使用collect_list将15亿条记录聚合为单个列表,导致DataFrame仅1个分区。写入时该分区的Task会在单个Executor上执行,需要加载整个聚合后的数据集到内存,16G的Executor内存无法容纳如此庞大的数据,触发OOM后被YARN NodeManager以外部信号终止(退出码143)。
    • 提交命令存在拼写错误:--cong应为--conf,且spark.dynamicAllocation.maxExecutor应为复数maxExecutors,导致动态资源配置未生效,无法按需扩展Executor资源。
  2. 多余空行问题:
    Spark默认JSON Writer会在每条记录后添加换行符,即使只有一条记录,也会在末尾生成空行,属于组件默认行为。

  3. 关于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转换为单个数组格式。

2. 解决多余空行问题

  • 方法一:写入后过滤空行
    • 本地文件:使用sed '/^$/d' input.json > output.json命令去除空行。
    • HDFS文件:先下载到本地处理后重新上传,或通过Spark读取文件并过滤空行后写入:
      spark.read.text("./TestJson").filter("value != ''").write.text("./TestJson_clean")
      
  • 方法二:自定义JSON写入逻辑
    将聚合后的数据拉到Driver,使用Python原生json模块生成无空行的JSON文件,完全控制输出格式:
    # 聚合后的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内存(如--driver-memory 32g)及spark.driver.maxResultSize(设为0表示无限制),确保Driver能容纳整个数据集。

内容的提问来源于stack exchange,提问作者Code Heaven

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 14:20:24