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

如何使用Spark Scala将字符串格式的复杂JSON转换为DataFrame

错误根因

你定义的schema要求Row包含Array[String]类型的info字段和String类型的more字段,但你传入Row的是完整的JSON字符串,Spark无法直接将整段字符串映射为你声明的复合结构,因此抛出类型不匹配异常。

实现方案

方案1:自动推导Schema(推荐,代码最简)

直接使用Spark原生的JSON读取能力解析字符串RDD,Spark会自动根据JSON结构推导Schema,无需手动定义,适配性最高。

import org.apache.spark.sql.SparkSession

// 初始化SparkSession
val ss = SparkSession.builder()
  .appName("JsonToDataFrame")
  .master("local[*]") // 本地测试用,生产环境移除该配置
  .getOrCreate()

// 你的原始JSON字符串
val data = """{
  "info": [
    {
      "done": "time",
      "id": 9,
      "type": "normal",
      "pid": 202020,
      "add": {
        "fields": true,
        "stat": "not sure"
      }
    },
    {
      "done": "time",
      "id": 14,
      "type": "normal",
      "pid": 764310,
      "add": {
        "fields": true,
        "stat": "sure"
      }
    },
    {
      "done": "time",
      "id": 9,
      "type": "normal",
      "pid": 202020,
      "add": {
        "note": {
          "id": 922,
          "score": 0
        }
      }
    }
  ],
  "more": {
    "a": "ok",
    "b": "fine",
    "c": 3
  }
}"""

// 构造字符串RDD后直接用json接口读取
val jsonRdd = ss.sparkContext.parallelize(Seq(data))
val df = ss.read.json(jsonRdd)

// 验证结果
df.printSchema()
df.show(false)

方案2:自定义Schema(适用于需要严格控制字段类型的场景)

如果需要避免自动推导的类型偏差,可以自定义和JSON结构完全匹配的Schema,通过from_json函数完成解析:

import org.apache.spark.sql.functions.{col, from_json}
import org.apache.spark.sql.types._

// 逐层定义和JSON结构匹配的Schema
val addSchema = new StructType()
  .add("fields", BooleanType)
  .add("stat", StringType)
  .add("note", new StructType()
    .add("id", LongType)
    .add("score", IntegerType)
  )

val infoItemSchema = new StructType()
  .add("done", StringType)
  .add("id", LongType)
  .add("type", StringType)
  .add("pid", LongType)
  .add("add", addSchema)

val rootSchema = new StructType()
  .add("info", ArrayType(infoItemSchema))
  .add("more", new StructType()
    .add("a", StringType)
    .add("b", StringType)
    .add("c", IntegerType)
  )

// 先将JSON字符串转为单列DataFrame,再解析
val rawDf = ss.createDataFrame(Seq((data,))).toDF("json_str")
val parsedDf = rawDf.select(from_json(col("json_str"), rootSchema).alias("root"))
val finalDf = parsedDf.select("root.*")

// 验证结果
finalDf.printSchema()
finalDf.show(false)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 15:57:02