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

PySpark3.x处理大JSON时repartition失效及写入ES报错问题

问题根因定位

1. 仅1个executor、单任务运行的根因

  • bz2属于不可分割压缩格式,Spark读取单个不可分割压缩文件时,读取阶段只能启动1个task加载全量数据,无法并行拆分读取。
  • 未显式指定JSON schema,Spark需要先启动单task扫描全量文件推断schema,该阶段也无法并行执行。
  • 本次任务在读取解析阶段就触发报错直接终止,还未执行到后续的repartition(20)逻辑,因此不会出现预期的20个并行任务。
  • spark-submit命令中错误添加了--class org.apache.spark.examples.SparkPi参数,该参数是Java/Scala Jar任务专用,Python脚本任务不需要配置,会干扰YARN的资源调度,导致executor分配异常。
  • 集群总核数仅40核,申请20个executor、每个2核共计40核,未预留系统及Hadoop基础服务占用的资源,会导致大量executor无法正常启动。

2. _corrupt_record报错的根因

读取JSON时未指定显式schema,当所有JSON行解析失败时,DataFrame只会保留默认的_corrupt_record错误列,Spark 2.3+版本禁止直接对仅包含该列的原始读取结果执行后续操作,即触发日志中抛出的AnalysisException。


解决方案

1. 解析报错修复

  • 提前定义显式JSON schema传入读取逻辑,避免Spark自动推断schema。可先抽取小部分样例数据解析后打印schema,再固化到代码中:
from pyspark.sql.types import StructType, StructField, StringType, LongType # 根据实际字段类型调整
# 自定义schema示例,替换为实际数据的字段结构
custom_schema = StructType([
    StructField("_id", StringType(), True),
    StructField("field1", LongType(), True),
    StructField("field2", StringType(), True)
])
df = spark.read.schema(custom_schema).option("multiline", "true").json(data_path)
  • 读取时添加容错配置并缓存结果,兼容解析失败的场景:
df = spark.read.option("multiline", "true") \
    .option("mode", "PERMISSIVE") \
    .option("columnNameOfCorruptRecord", "_corrupt_record") \
    .schema(custom_schema) \
    .json(data_path).cache()
# 过滤掉解析失败的行
df = df.filter(df["_corrupt_record"].isNull()).drop("_corrupt_record")
  • 校验源JSON文件格式合法性,确认多行JSON的顶层对象结构完整、无语法错误。

2. 并行度&资源配置优化

  • 拆分原2GB的bz2文件为10~20个大小相近的小bz2文件,Spark支持对多个不可分割压缩文件并行读取,每个文件对应1个读取task,从读取阶段就实现并行执行。
  • 修正spark-submit命令,移除多余的class参数,调整资源配置贴合集群实际容量:
spark-submit --jars elasticsearch-spark-30_2.12-7.14.1.jar \
    --master yarn \
    --deploy-mode cluster \
    --driver-memory 8g \
    --executor-memory 6g \
    --num-executors 8 \
    --executor-cores 4 \
    myscript.py
  • 解析后的DataFrame执行persist()缓存后再调用repartition(20),确保重分区逻辑生效,后续写入ES阶段可并行执行。

内容的提问来源于stack exchange,提问作者Hafiz Muhammad Shafiq

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 15:48:04