Spark Streaming报错:不支持org.apache.spark.sql.types.DataType类型Schema如何解决
问题根因
- Spark UDF的返回值必须是Spark SQL支持的可序列化存储类型,
org.apache.spark.sql.types.DataType属于Spark内部的schema描述类,没有对应SQL存储格式,因此直接注册返回该类型的UDF必然抛出不支持的异常。 - 内置
from_json函数的第二个schema参数要求是固定的常量值,不支持传入随行变化的Column类型,因此无法直接通过内置函数实现按行匹配专属schema解析的需求。
适配流处理的解决方案
直接封装一个接收JSON字符串、Schema DDL字符串两个参数的UDF,在UDF内部完成DDL转Schema、JSON解析的全流程,避开内置函数的限制:
import org.apache.spark.sql.types.{DataType, StructType} import org.apache.spark.sql.functions.udf import org.apache.spark.sql.catalyst.json.JacksonParser // 注册自定义解析UDF val parseJsonWithDynamicSchema = udf((jsonStr: String, schemaDdl: String) => { val targetSchema = DataType.fromDDL(schemaDdl).asInstanceOf[StructType] // 调用Spark内置JSON解析器按指定schema解析 JacksonParser.parseJson(jsonStr, targetSchema, Map.empty[String, String]).headOption.orNull }) // 逐行解析生成结果表 val out = input.select( parseJsonWithDynamicSchema(col("value").cast("string"), col("schema")).alias("parsed_json") )
额外说明
- 该方案不需要提前收集全量schema,也不依赖固定schema,完全适配流处理场景
- 如果同一微批次内所有行的解析后结构一致,可以直接用
out.select("parsed_json.*")把结构体字段展开为独立列 - 如果不同行的schema结构差异极大,建议保留结构体格式存储,避免展开后出现大量空值列
内容的提问来源于stack exchange,提问作者Khan Saab
相关产品推荐
相关产品推荐

