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

含Union类型的Avro对象与JSON互转及Kafka发送问题

Quarkus Outbox模式下Avro Union类型JSON序列化问题解决

问题根源

你的details字段是Avro对象类型的Union(OtherUpdateDetails/SpecialUpdateDetails),默认Jackson序列化时不会写入类型标识,从数据库读回反序列化时,Jackson无法判断该映射到Union里的哪个具体类型,导致生成的Avro对象结构不符合Kafka Producer要求,最终触发报错。

解决步骤

1. 补全依赖

确保pom.xml中包含以下依赖(版本需匹配你的Quarkus和Avro版本):

<dependency>
    <groupId>org.apache.avro</groupId>
    <artifactId>avro</artifactId>
    <version>${avro.version}</version>
</dependency>
<dependency>
    <groupId>com.fasterxml.jackson.dataformat</groupId>
    <artifactId>jackson-dataformat-avro</artifactId>
    <version>${jackson.version}</version>
</dependency>
<dependency>
    <groupId>io.quarkus</groupId>
    <artifactId>quarkus-jackson</artifactId>
</dependency>

2. 自定义Jackson Mapper保留Union类型信息

默认Jackson不会输出Avro Union的类型标识,需自定义ObjectMapper开启该功能:

import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.dataformat.avro.AvroMapper;
import jakarta.enterprise.context.ApplicationScoped;
import jakarta.enterprise.inject.Produces;

@ApplicationScoped
public class AvroJacksonConfig {

    @Produces
    public ObjectMapper avroObjectMapper() {
        AvroMapper avroMapper = new AvroMapper();
        // 强制写入Union类型的标识信息
        avroMapper.enable(AvroMapper.Feature.WRITE_UNION_TYPE_ID);
        return avroMapper;
    }
}

3. 序列化Avro对象到数据库

使用自定义的AvroMapper序列化,确保JSON中携带类型标识:

@Inject
ObjectMapper avroObjectMapper;

public String serializeMessage(ValidatedUpdate message) throws IOException {
    return avroObjectMapper.writerWithSchemaFor(ValidatedUpdate.class).writeValueAsString(message);
}

4. 从数据库反序列化Avro对象

反序列化时指定目标Schema,保证类型正确映射:

public ValidatedUpdate deserializeMessage(String json) throws IOException {
    return avroObjectMapper.readerFor(ValidatedUpdate.class).with(ValidatedUpdate.getClassSchema()).readValue(json);
}

5. Kafka Producer配置

确保使用Avro序列化器,在Quarkus的application.properties中添加:

quarkus.kafka.producer.value.serializer=io.confluent.kafka.serializers.KafkaAvroSerializer
quarkus.kafka.producer.schema.registry.url=http://你的schema-registry地址:8081

验证要点

  • 检查数据库中的JSON字符串,details字段需包含type标识,示例:
    "details": {
      "type": "com.acme.kafka.SpecialUpdateDetails",
      "stuff": "demo",
      "isYes": true
    }
    
  • 若仍报错,抓取具体错误信息,排查是否存在Schema不匹配或类型标识缺失的问题。

内容的提问来源于stack exchange,提问作者Nan0fire

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 03:22:57