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

Spark如何直接读取单行多JSON对象文件并生成多行DataFrame

读取单行多JSON对象文件的可行方案

原实现问题说明

你现有代码使用sc.textFile(path).collect()(0)会将整个文件内容全部拉取到Driver节点内存,仅适合测试小体积文件,生产环境遇到大文件会直接触发Driver OOM,不具备可扩展性。

方案1:分布式流式处理(全Spark版本支持,推荐)

直接在Executor端做字符串拆分和补全,全程数据不落Driver,支持任意大小的文件:

import spark.implicits._
import org.apache.spark.sql.functions._

// 读取为仅含单行文本的DataFrame
val rawTextDF = spark.read.text(path)

// 分布式拆分并补全每个独立JSON对象
val jsonStrDS = rawTextDF.flatMap(row => {
  val fullLine = row.getString(0)
  fullLine.split("\\}\\,\\{")
    .map(part => {
      val withHead = if (!part.startsWith("{")) s"{$part" else part
      if (!withHead.endsWith("}")) s"$withHead}" else withHead
    })
})

// 直接解析JSON字符串数据集得到结构化结果
val resultDF = spark.read.json(jsonStrDS)

方案2:利用行分隔符配置简化逻辑(Spark 3.0+支持)

通过设置自定义行分隔符直接拆分文件,再补全JSON格式即可:

import org.apache.spark.sql.functions._

val jsonStrDS = spark.read
  // 设置行分隔符为`},`,直接拆分出每个JSON对象主体
  .option("lineSep", "},")
  .text(path)
  // 补全每个JSON对象缺失的右括号
  .select(
    when(col("value").endsWith("}"), col("value"))
      .otherwise(concat(col("value"), lit("}")))
      .alias("value")
  )
  .as[String]

// 解析JSON得到结构化数据
val resultDF = spark.read.json(jsonStrDS)

注意事项

如果JSON字段的value中包含},{或},的子串,上述拆分逻辑会报错,这种场景建议先将源文件预处理为标准JSON Lines格式(每个JSON对象独占一行),再用原生spark.read.json(path)读取,是兼容性最高的方案。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 11:06:00