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

使用KafkaListener反序列化Avro原始字符串消息键失败求助

解决Kafka Avro类型消息键的反序列化问题

问题原因

你的消息键是用Avro原始string类型定义的,Avro序列化这类数据时,会在字符串字节前添加变长长度前缀(采用zig-zag编码)。直接用StringDeserializer反序列化时,会把这些前缀字节当成普通字符串内容解析,所以得到了带\n&的异常结果。

正确解决方案

需要使用Avro专用的反序列化器KafkaAvroDeserializer来处理,具体配置如下:

1. 修改@KafkaListener配置

在监听器的properties中添加Avro反序列化相关配置:

@KafkaListener(
    topics = {"topic"},
    autoStartup = "true",
    properties = {
        // 指定Avro键反序列化器
        ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG + "=io.confluent.kafka.serializers.KafkaAvroDeserializer",
        // 替换为你的Schema Registry地址
        "schema.registry.url=http://localhost:8081",
        // 因为是Avro原始类型而非自定义SpecificRecord,所以设为false
        "specific.avro.reader=false"
    }
)
@Transactional
public void consume(@Payload(required = false) Node skill, @Header(OFFSET) Long offset,
                    @Header(RECEIVED_MESSAGE_KEY) String messageKey, @Header(RECEIVED_TOPIC) String topic,
                    @Header(RECEIVED_TIMESTAMP) Long timestamp) {
    // 此时messageKey会是预期的"7770000000000105411"
}

2. 确保依赖正确

需要引入Confluent的Avro序列化依赖(版本请匹配你的Kafka和Schema Registry版本):

<!-- Maven依赖示例 -->
<dependency>
    <groupId>io.confluent</groupId>
    <artifactId>kafka-avro-serializer</artifactId>
    <version>7.4.0</version>
</dependency>

额外说明

如果你的项目中已经全局配置了消费者工厂,也可以在工厂层面统一设置这些参数,避免每个监听器重复配置:

@Bean
public ConsumerFactory<String, Node> consumerFactory() {
    Map<String, Object> configProps = new HashMap<>();
    configProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
    configProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, KafkaAvroDeserializer.class);
    configProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class); // 假设值用JSON反序列化
    configProps.put("schema.registry.url", "http://localhost:8081");
    configProps.put("specific.avro.reader", false);
    return new DefaultKafkaConsumerFactory<>(configProps);
}

内容的提问来源于stack exchange,提问作者Никита Спиридонов

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 08:40:24