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

