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

Scala将DataFrame的JSON列转为数组写入HDFS时遇容器失败问题

问题分析与解决方案

原代码失败原因

你的代码中使用groupBy("id").agg(collect_list("value"))会将每个分区内的2M条JSON字符串全部加载到内存中生成一个超大列表,这直接导致Executor内存溢出,触发YARN容器被强制杀死(退出码143对应外部信号终止),最终抛出InterruptedException和SparkException。

优化方案

核心思路是避免将整个分区的数据一次性加载到内存,改为流式处理每个分区的记录,直接生成符合要求的JSON数组格式文件。具体实现如下:

步骤1:计算目标文件数量

根据50M总记录和每个文件2M记录的要求,计算分区数:

val targetFileCount = 50000000 / 2000000 // 结果为25

步骤2:流式处理分区生成JSON数组

使用mapPartitions迭代处理每个分区的记录,逐行生成JSON数组的内容,无需一次性加载所有数据到内存:

import org.apache.spark.sql.SaveMode

Targetdf.repartition(targetFileCount)
  .mapPartitions { iter =>
    if (!iter.hasNext) {
      // 空分区直接返回空迭代器
      Iterator.empty
    } else {
      // 处理第一个记录,生成数组开头
      val firstJson = iter.next().getAs[String]("value")
      val start = Iterator("[", firstJson)
      
      // 处理剩余记录,每个记录前添加逗号
      val middle = iter.map(row => s", ${row.getAs[String]("value")}")
      
      // 添加数组结尾
      val end = Iterator("]")
      
      // 拼接所有部分
      start ++ middle ++ end
    }
  }
  .write
  .mode(SaveMode.Overwrite)
  .text(HdfsPath)

代码说明

  1. repartition(targetFileCount):将数据均匀分配到25个分区,确保每个分区对应一个输出文件,且每个文件包含约2M条记录。
  2. mapPartitions流式处理:
    • 不将整个分区的数据转换为列表,而是逐个迭代记录
    • 先输出数组开头[和第一条JSON字符串
    • 后续每条记录前添加, ,避免末尾多余逗号
    • 最后输出数组结尾]
  3. text格式写入:直接输出文本内容,每个分区生成一个文件,文件内容为完整的JSON数组(允许换行,JSON语法支持任意空白字符)。

额外优化建议

  • 调整Executor内存:如果单条JSON字符串过大(比如超过1KB),可以适当提高Executor内存配置(如--executor-memory 4G),避免处理过程中出现内存不足。
  • 检查数据分布:如果原数据存在分区倾斜,可考虑使用repartitionByRange或自定义分区键优化数据分布。
  • 验证输出格式:可以读取单个输出文件,确认其为合法的JSON数组格式。

内容的提问来源于stack exchange,提问作者Amiya Ghosh

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 15:07:06