为什么Spark中schema_of_json传入列对象时会执行失败?
问题原因
schema_of_json是Spark设计的静态Schema生成函数,入参要求是查询计划生成阶段就能确定的字符串字面量,不支持运行时动态计算的列值,底层逻辑和Spark的设计机制直接相关:
- Spark DataFrame的Schema是强类型、静态绑定的,所有列的类型必须在作业启动前的逻辑计划、物理计划生成阶段完全确定,不能在运行时随数据动态变化。如果传入列作为参数,Spark无法在计划阶段预知列内存储的JSON内容,自然没法确定输出的
schemaDetected列的类型,所以直接报参数类型不匹配错误。 - 第一种写法直接把JSON字符串写在
lit里传入函数,Spark在计划阶段就能直接拿到字符串内容,当场解析出对应Schema,所以可以正常运行。第二种写法把JSON存到了literal列中,哪怕所有行的内容完全一致,Spark在计划阶段也不会读取具体数据值推导参数,只会识别为动态列,不符合入参要求。
解决方案
根据你的业务场景可以选择两种处理方式:
场景1:列中所有JSON结构一致
这是最常见的场景,先采样出一条JSON样本,提前解析出Schema再应用到全量数据即可:
// 1. 取一条JSON样本数据 val sampleJsonStr = df.select("literal").head().getString(0) // 2. 提前生成对应的JSON Schema val jsonSchema = schema_of_json(lit(sampleJsonStr)) // 3. 用预生成的Schema处理全量数据 val resultDf = df .withColumn("schemaDetected", lit(jsonSchema.simpleString)) .withColumn("parsedJson", from_json(col("literal"), jsonSchema)) resultDf.show(false)
场景2:列中每行JSON结构不同
如果每行的JSON结构都不一样,需要逐行返回对应Schema的字符串,可以通过自定义UDF实现,注意该方式需要逐行解析JSON,数据量大时会有明显性能损耗:
import org.apache.spark.sql.functions.udf import com.fasterxml.jackson.databind.ObjectMapper import com.fasterxml.jackson.module.scala.DefaultScalaModule import org.apache.spark.sql.types.DataType val jsonSchemaUdf = udf((jsonStr: String) => { if (jsonStr == null || jsonStr.trim.isEmpty) { null } else { val mapper = new ObjectMapper().registerModule(DefaultScalaModule) // 解析JSON后推导Schema返回字符串格式 DataType.fromJson(schema_of_json(lit(jsonStr)).expr.eval().toString).simpleString } }) val resultDf = df.withColumn("schemaDetected", jsonSchemaUdf(col("literal"))) resultDf.show(false)
内容的提问来源于stack exchange,提问作者MrElephant
相关产品推荐
相关产品推荐

