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

Scala中使用from_json()解析DataFrame多行JSON数据的问题

问题解析与解决方法

你的问题核心在于Schema定义不匹配JSON列的实际结构,以及缺少数组展开的关键步骤:

1. 为什么原代码只返回单条记录?

你定义的schema是单个JSON对象的结构,但你的JSON列存储的是JSON数组。当你用from_json(col("JSON"), schema)时,Spark会尝试把整个数组当成单个对象解析,最终只会提取数组的第一个元素,所以只得到一行结果。

2. 正确的解决步骤

第一步:修正Schema,适配JSON数组结构

需要把你定义的单个对象Schema包裹在ArrayType中,明确告诉Spark这是一个数组类型的JSON:

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

// 单个JSON对象的Schema
val singleObjSchema = StructType(Array(
  StructField("Hour", IntegerType),
  StructField("Total", IntegerType),
  StructField("Fail", IntegerType)
))

// 适配JSON数组的Schema
val jsonArraySchema = ArrayType(singleObjSchema)

第二步:解析JSON数组并展开为多行

先用from_json解析数组,再用explode函数把数组中的每个元素拆成单独的行:

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

val newDF = DF
  // 解析JSON数组为Spark数组类型列
  .withColumn("parsed_json_array", from_json(col("JSON"), jsonArraySchema))
  // 展开数组,每个元素对应一行
  .select(explode(col("parsed_json_array")).alias("json_obj"))
  // 提取对象中的字段
  .select(col("json_obj.*"))

newDF.show()

3. 简化写法(可选)

你也可以把Schema的定义直接内嵌到代码中,减少变量定义:

val newDF = DF
  .withColumn("parsed_json", from_json(col("JSON"), ArrayType(StructType(Array(
    StructField("Hour", IntegerType),
    StructField("Total", IntegerType),
    StructField("Fail", IntegerType)
  )))))
  .select(explode(col("parsed_json")).alias("json_obj"))
  .select("json_obj.*")

newDF.show()

执行以上代码后,就能得到你期望的三行输出了。

内容的提问来源于stack exchange,提问作者S. Rud

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 09:56:50