Spark DataFrame借助Schema Registry的Avro写入HDFS报错原因排查
我之前也碰到过一模一样的问题,咱们一步步拆解来看看:
问题重现
你尝试将带有timestamp类型且允许为空的CreateDate字段的Spark DataFrame,通过Schema Registry获取的Avro Schema写入HDFS时,触发了如下报错:
Caused by: org.apache.avro.AvroRuntimeException: Not a union: {"type":"long","logicalType":"timestamp-millis"}
你的Avro Schema定义的该字段是包含null的union类型:
{"name":"CreateDate","type":["null",{"type":"long","logicalType":"timestamp-millis"}],"default":null}
而DataFrame里的字段类型是:
|-- CreateDate: timestamp (nullable = true)
写入代码大致是:
dataDF.write .mode("append") .format("avro") .option( "avroSchema", SchemaRegistry.getSchema( schemaRegistryConfig.url, schemaRegistryConfig.dataSchemaSubject, schemaRegistryConfig.dataSchemaVersion)) .save(hdfsURL)
报错原因
核心问题在于Spark Avro序列化器期望nullable的DataFrame字段对应Avro的union类型(必须包含null),但实际从Schema Registry获取到的CreateDate字段类型并非union,而是单独的{"type":"long","logicalType":"timestamp-millis"}。
具体来说,Spark的AvroSerializer里的resolveNullableType方法会做校验:如果DataFrame字段是nullable = true,对应的Avro类型必须是包含null的union。一旦发现拿到的Avro类型是单个非union类型,就会抛出这个"Not a union"的异常。
为什么会拿到非union的schema?大概率是这两个原因之一:
- Schema Registry中存储的对应subject、版本的schema本身就不是你预期的union类型(可能是之前错误上传的);
- 调用
SchemaRegistry.getSchema(...)的逻辑有问题,比如序列化/反序列化schema时丢失了union结构,或者传错了subject/版本号导致拿到了错误的schema。
解决方案
先验证Schema Registry中的schema正确性
直接通过Schema Registry的API或者管理UI,查询对应dataSchemaSubject和版本的schema,确认CreateDate字段确实是包含null的union类型。如果发现不对,重新上传正确的schema到Registry。检查schema获取代码的正确性
在写入代码前,先把SchemaRegistry.getSchema(...)返回的schema打印出来,对比是否和你预期的一致。比如加一行:val avroSchemaStr = SchemaRegistry.getSchema(...) println(avroSchemaStr)如果打印出来的schema里
CreateDate不是union,那就得排查获取逻辑的问题——比如是否有代码意外修改了schema字符串,或者subject/版本参数传错了。检查Spark Avro版本兼容性
某些旧版本的Spark Avro模块(比如Spark 3.0及更早)在处理带logicalType的union类型时可能存在bug,尝试升级到Spark 3.1+对应的Avro依赖版本,看看是否能解决问题。临时应急方案(不推荐长期使用)
如果暂时无法修改Registry的schema,可以先把DataFrame的CreateDate字段转换为nullable的long类型(对应timestamp-millis的格式):import org.apache.spark.sql.functions._ val modifiedDF = dataDF.withColumn("CreateDate", when(col("CreateDate").isNotNull, unix_timestamp(col("CreateDate")) * 1000).otherwise(null))再用修改后的DataFrame写入,不过这只是权宜之计,最好还是保证Schema和数据类型的匹配。
内容的提问来源于stack exchange,提问作者Cassie

