You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark DataFrame借助Schema Registry的Avro写入HDFS报错原因排查

解决Spark Avro写入时抛出"Not a union"的报错问题

我之前也碰到过一模一样的问题,咱们一步步拆解来看看:

问题重现

你尝试将带有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。

解决方案

  1. 先验证Schema Registry中的schema正确性
    直接通过Schema Registry的API或者管理UI,查询对应dataSchemaSubject和版本的schema,确认CreateDate字段确实是包含null的union类型。如果发现不对,重新上传正确的schema到Registry。

  2. 检查schema获取代码的正确性
    在写入代码前,先把SchemaRegistry.getSchema(...)返回的schema打印出来,对比是否和你预期的一致。比如加一行:

    val avroSchemaStr = SchemaRegistry.getSchema(...)
    println(avroSchemaStr)
    

    如果打印出来的schema里CreateDate不是union,那就得排查获取逻辑的问题——比如是否有代码意外修改了schema字符串,或者subject/版本参数传错了。

  3. 检查Spark Avro版本兼容性
    某些旧版本的Spark Avro模块(比如Spark 3.0及更早)在处理带logicalType的union类型时可能存在bug,尝试升级到Spark 3.1+对应的Avro依赖版本,看看是否能解决问题。

  4. 临时应急方案(不推荐长期使用)
    如果暂时无法修改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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.13 08:24:29