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

如何在使用@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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.23 21:24:00