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
相关产品推荐
相关产品推荐

