使用Spark from_json转换为Dataset时出现列无法解析异常
问题原因分析
你遇到的AnalysisException核心原因是转换Dataset时,Spark无法将当前DataFrame的列结构与MyMeta类的字段匹配:
执行val df5 = df4.withColumn("parsedMeta", from_json(col("meta"), metaSchema)).drop("meta").as[MyMeta]后,df5的结构仅包含一个parsedMeta列(类型为对应MyMeta的结构体),但MyMeta类需要的是op和table两个独立顶级字段,Spark找不到这两个字段,因此抛出解析错误。
解决方案
以下两种方式可快速修复该问题:
方式1:展开结构体字段后转换Dataset
先将parsedMeta中的子字段提取为DataFrame的顶级列,再映射到MyMeta类:
val df5 = df4.withColumn("parsedMeta", from_json(col("meta"), metaSchema)) .select("parsedMeta.op", "parsedMeta.table") .as[MyMeta]
方式2:直接从结构体列映射Dataset
利用Spark结构体列与case class的直接映射能力,无需展开字段,只需确保列结构匹配:
val df5 = df4.withColumn("parsedMeta", from_json(col("meta"), metaSchema)) .select("parsedMeta") .as[MyMeta] // 更简洁的写法 val df5 = df4.select(from_json(col("meta"), metaSchema).as("parsedMeta")) .select("parsedMeta.*") .as[MyMeta]
额外优化建议
你最初的DataFrame处理步骤可大幅简化,无需手动处理字符串转Map再提取元素的流程,直接用Spark的JSON读取器解析即可:
case class MyMeta(op: Option[String], table: Option[String]) case class InputData(meta: MyMeta, data: Map[String, String], key: Map[String, String]) val inputSchema = Encoders.product[InputData].schema val path = "/FileStore/tables/json_0002C_file.txt" val df = spark.read.json(path) // 直接读取JSON格式文件,自动解析结构 val metaDs = df.select("meta").as[MyMeta] metaDs.show(false) metaDs.printSchema()
内容的提问来源于stack exchange,提问作者Ged
相关产品推荐
相关产品推荐

