使用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,提问作者Никита Спиридонов
相关产品推荐
相关产品推荐

