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

Spark Scala中如何从DataFrame的JSON列提取嵌套字段值?

问题:从DataFrame的JSON列中嵌套提取字段值并断言

测试代码片段

class MyTest extends AnyFlatSpec with Matchers {

   ....

  it should "calculate" in {

    val testDf= Seq(
      testDf(1, "customer1", "Hi"),
      testDf(1, "customer2", "Hi")
    ).toDS().toDF()


    val out = MyClass.procOut(spark, testDf)
    out.count() should be(1)
    // 此处写法有问题,无法正确提取JSON嵌套字段
    out.where(col("customer_id")==="customer1").first().getString(output.first().fieldIndex("json_col")) should be(?) 

  }

}

问题说明

out是一个DataFrame,其中json_col列存储的JSON结构如下:

{
  "statistics": {
    "Group2": {
      "buy": 1
    }
  }
}

需要提取buy字段的值并断言其为1,之前尝试的写法语法错误,求正确的提取方式。


解决方案

方法1:使用Spark内置JSON函数(推荐)

利用Spark的get_json_object函数直接在DataFrame层面解析JSON,无需提取字符串后再处理:

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

// 提取buy字段并转成整数类型
val buyValue = out.where(col("customer_id") === "customer1")
  .select(get_json_object(col("json_col"), "$.statistics.Group2.buy").cast(IntegerType))
  .first()
  .getInt(0)

// 断言值为1
buyValue should be(1)

说明:get_json_object的第二个参数是JSONPath表达式,$.statistics.Group2.buy表示从根节点开始逐层定位到buy字段。

方法2:提取JSON字符串后用Scala JSON库解析

如果需要手动解析JSON字符串,可以使用Jackson(Spark默认依赖):

import com.fasterxml.jackson.databind.ObjectMapper

// 先提取json_col的字符串内容
val jsonStr = out.where(col("customer_id") === "customer1")
  .first()
  .getString(out.schema.fieldIndex("json_col"))

// 用Jackson解析JSON并提取buy值
val mapper = new ObjectMapper()
val jsonNode = mapper.readTree(jsonStr)
val buyValue = jsonNode.get("statistics").get("Group2").get("buy").asInt()

// 断言
buyValue should be(1)

错误原因说明

你之前的写法错误地将**列索引(整数类型)**和JSON字段索引混在一起,fieldIndex("json_col")返回的是列在DataFrame中的位置序号,不能直接用["statistics"]这种JSON字段访问方式操作,因此导致语法错误。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 02:40:33