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

如何配置错误处理器丢弃反序列化错误且不投递至死信主题?

问题描述

我尝试按如下方式配置自定义错误处理器,以实现主题中的无效消息不重试且不投递至死信主题:

public ConcurrentKafkaListenerContainerFactory<K, V> kafkaListenerContainerFactory(
      Map<String, Object> consumerProps) {
  factory.setConsumerFactory(consumerFactory(consumerProps));
  
  final var eh = new DefaultErrorHandler(new FixedBackOff(0L, 0L));
  eh.setAckAfterHandle(true);
  eh.setResetStateOnRecoveryFailure(false);
  eh.setLogLevel(KafkaException.Level.DEBUG);
  eh.addNotRetryableExceptions(SerializationException.class);
  eh.addNotRetryableExceptions(RecordDeserializationException.class);
  eh.addNotRetryableExceptions(IllegalArgumentException.class);

  factory.setCommonErrorHandler(eh);

  return factory;
}

但问题在于,RecordDeserializationException未被错误处理器捕获,报错信息如下:

2025-03-11 21:39:09.375 - - ERROR o.s.k.l.KafkaMessageListenerContainer service-reference-java-heartRate-0-C-1 - Consumer exception 
java.lang.IllegalStateException: This error handler cannot process 'SerializationException's directly; please consider configuring an 'ErrorHandlingDeserializer' in the value and/or key deserializer
        at org.springframework.kafka.listener.DefaultErrorHandler.handleOtherException(DefaultErrorHandler.java:192)
        at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.handleConsumerException(KafkaMessageListenerContainer.java:1934)
        at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:1365)
        at java.base/java.util.concurrent.CompletableFuture$AsyncRun.run(CompletableFuture.java:1804)
        at java.base/java.lang.Thread.run(Thread.java:833)
Caused by: org.apache.kafka.common.errors.RecordDeserializationException: Error deserializing key/value for partition com........readings2-0 at offset 0. If needed, please seek past the record to continue consumption.
        at org.apache.kafka.clients.consumer.internals.CompletedFetch.parseRecord(CompletedFetch.java:309)
        at org.apache.kafka.clients.consumer.internals.CompletedFetch.fetchRecords(CompletedFetch.java:263)
        at org.apache.kafka.clients.consumer.internals.AbstractFetch.fetchRecords(AbstractFetch.java:340)
        at org.apache.kafka.clients.consumer.internals.AbstractFetch.collectFetch(AbstractFetch.java:306)
        at org.apache.kafka.clients.consumer.KafkaConsumer.pollForFetches(KafkaConsumer.java:1235)
        at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1186)
        at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1159)
        at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollConsumer(KafkaMessageListenerContainer.java:1649)
        at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doPoll(KafkaMessageListenerContainer.java:1624)
        at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollAndInvoke(KafkaMessageListenerContainer.java:1421)
        at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:1313)
        ... 2 common frames omitted
Caused by: java.lang.IllegalArgumentException: com.google.protobuf.InvalidProtocolBufferException: Type of the Any message does not match the given class.

请问我哪里配置有误?如何确保RecordDeserializationException被错误处理器捕获且不进行重试?


解决方案

核心问题是未配置ErrorHandlingDeserializer,Spring Kafka的DefaultErrorHandler无法直接处理消费端反序列化阶段抛出的RecordDeserializationException——这类异常发生在Kafka客户端拉取消息后、Spring监听器处理消息前,必须通过ErrorHandlingDeserializer将异常封装,才能传递给DefaultErrorHandler处理。

具体修复步骤:

  1. 用ErrorHandlingDeserializer包装实际反序列化器
    修改消费者配置,将原有的key/value反序列化器替换为ErrorHandlingDeserializer,并指定实际的反序列化器类:

    // 配置value的ErrorHandlingDeserializer,替换为你实际使用的反序列化器
    consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);
    consumerProps.put(ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS, YourActualProtobufDeserializer.class.getName());
    
    // 如果key也需要处理,同样配置key的ErrorHandlingDeserializer
    consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);
    consumerProps.put(ErrorHandlingDeserializer.KEY_DESERIALIZER_CLASS, YourActualKeyDeserializer.class.getName());
    
  2. 保留原有DefaultErrorHandler配置
    你当前的DefaultErrorHandler已正确设置了不重试(FixedBackOff(0L, 0L))、标记非重试异常,无需修改。配置ErrorHandlingDeserializer后,反序列化异常会被包装成ListenerExecutionFailedException并携带原始异常,此时DefaultErrorHandler就能识别并按规则处理。

  3. 可选:自定义反序列化异常处理逻辑
    如果需要对反序列化失败的消息做特殊处理(比如记录原始消息内容),可以给ErrorHandlingDeserializer设置自定义DeserializationExceptionHandler,默认逻辑已满足将异常传递给DefaultErrorHandler的需求。

原配置不生效的原因

RecordDeserializationException是在Kafka客户端拉取线程中抛出的,未进入Spring Kafka的监听器执行链路,DefaultErrorHandler只能处理监听器执行阶段的异常。而ErrorHandlingDeserializer作为代理反序列化器,会捕获反序列化异常,将消息封装成包含异常的ConsumerRecord后传递给监听器,此时异常才能被DefaultErrorHandler捕获处理。


内容的提问来源于stack exchange,提问作者Higher-Kinded Type

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 22:33:17