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)
代码说明
repartition(targetFileCount):将数据均匀分配到25个分区,确保每个分区对应一个输出文件,且每个文件包含约2M条记录。mapPartitions流式处理:- 不将整个分区的数据转换为列表,而是逐个迭代记录
- 先输出数组开头
[和第一条JSON字符串 - 后续每条记录前添加
,,避免末尾多余逗号 - 最后输出数组结尾
]
text格式写入:直接输出文本内容,每个分区生成一个文件,文件内容为完整的JSON数组(允许换行,JSON语法支持任意空白字符)。
额外优化建议
- 调整Executor内存:如果单条JSON字符串过大(比如超过1KB),可以适当提高Executor内存配置(如
--executor-memory 4G),避免处理过程中出现内存不足。 - 检查数据分布:如果原数据存在分区倾斜,可考虑使用
repartitionByRange或自定义分区键优化数据分布。 - 验证输出格式:可以读取单个输出文件,确认其为合法的JSON数组格式。
内容的提问来源于stack exchange,提问作者Amiya Ghosh
相关产品推荐
相关产品推荐

