如何在使用@KafkaListener的Kafka消费者中记录不兼容Avro schema的消息
实现方案
Avro schema兼容性校验失败的异常会在消息进入@KafkaListener标注的业务方法之前抛出,你可以通过以下两种成熟方案实现异常日志记录:
方案1:用ErrorHandlingDeserializer包装Avro反序列化器
这是Spring Kafka官方推荐的标准方案,通过自带的反序列化异常包装器捕获所有反序列化阶段异常,再统一处理日志:
- 首先修改消费者配置,把原有Avro反序列化器配置为
ErrorHandlingDeserializer的委托类:
spring: kafka: consumer: value-deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer properties: # 配置实际的Avro反序列化器作为委托类 spring.deserializer.value.delegate.class: io.confluent.kafka.serializers.KafkaAvroDeserializer # 保留原有Avro相关配置 schema.registry.url: http://你的schema-registry地址:8081 specific.avro.reader: true # 如果需要同时处理key的反序列化异常,同步添加以下配置即可 # key-deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer # spring.deserializer.key.delegate.class: io.confluent.kafka.serializers.KafkaAvroDeserializer
- 配置
ConcurrentKafkaListenerContainerFactory添加异常处理器,捕获Avro序列化异常记录日志:
@Bean public ConcurrentKafkaListenerContainerFactory<?, ?> kafkaListenerContainerFactory( ConcurrentKafkaListenerContainerFactoryConfigurer configurer, ConsumerFactory<Object, Object> consumerFactory) { ConcurrentKafkaListenerContainerFactory<Object, Object> factory = new ConcurrentKafkaListenerContainerFactory<>(); configurer.configure(factory, consumerFactory); // 自定义异常处理器 factory.setCommonErrorHandler(new CommonErrorHandler() { @Override public void handleOtherException(Exception thrownException, Consumer<?, ?> consumer, MessageListenerContainer container, boolean batchListener) { // 匹配Avro schema兼容性相关异常 if (thrownException instanceof org.apache.kafka.common.errors.SerializationException || thrownException instanceof io.confluent.kafka.schemaregistry.client.rest.exceptions.RestClientException) { // 按需记录消息元信息:topic、分区、偏移量、异常栈等,方便定位脏数据 log.error("Avro schema兼容性校验失败,消息无法正常消费", thrownException); // 可根据业务需求选择提交偏移量,避免单条脏数据阻塞整个消费组 consumer.commitAsync(); } } }); return factory; }
方案2:自定义Avro反序列化器直接捕获异常
如果不想依赖Spring的ErrorHandlingDeserializer,可以直接继承官方KafkaAvroDeserializer,在反序列化逻辑中嵌入日志记录逻辑:
public class LoggableKafkaAvroDeserializer extends KafkaAvroDeserializer { private static final Logger log = LoggerFactory.getLogger(LoggableKafkaAvroDeserializer.class); @Override public Object deserialize(String topic, byte[] data) { try { return super.deserialize(topic, data); } catch (SerializationException e) { log.error("Topic {} 存在不兼容Avro schema的消息,原始字节长度:{}", topic, data == null ? 0 : data.length, e); // 按需选择抛出异常走原有错误流,或者返回null传递给业务代码处理 throw e; } } @Override public Object deserialize(String topic, Headers headers, byte[] data) { try { return super.deserialize(topic, headers, data); } catch (SerializationException e) { log.error("Topic {} 存在不兼容Avro schema的消息,消息头:{},原始字节长度:{}", topic, headers, data == null ? 0 : data.length, e); throw e; } } }
配置完成后把消费者的反序列化器类替换为你自定义的LoggableKafkaAvroDeserializer即可生效。
注意事项
- 批量消费场景需要对应调整错误处理器的批量异常处理逻辑,避免单条异常影响整批消息消费
- 日志建议记录消息的topic、分区、偏移量、原始字节、关联schema id等信息,方便后续定位脏数据来源
- 可根据业务需求配置死信队列,把不兼容的消息转发至死信队列存储,避免直接丢弃
内容的提问来源于stack exchange,提问作者Marinos An
相关产品推荐
相关产品推荐

