如何配置错误处理器丢弃反序列化错误且不投递至死信主题?
我尝试按如下方式配置自定义错误处理器,以实现主题中的无效消息不重试且不投递至死信主题:
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处理。
具体修复步骤:
用
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());保留原有
DefaultErrorHandler配置
你当前的DefaultErrorHandler已正确设置了不重试(FixedBackOff(0L, 0L))、标记非重试异常,无需修改。配置ErrorHandlingDeserializer后,反序列化异常会被包装成ListenerExecutionFailedException并携带原始异常,此时DefaultErrorHandler就能识别并按规则处理。可选:自定义反序列化异常处理逻辑
如果需要对反序列化失败的消息做特殊处理(比如记录原始消息内容),可以给ErrorHandlingDeserializer设置自定义DeserializationExceptionHandler,默认逻辑已满足将异常传递给DefaultErrorHandler的需求。
原配置不生效的原因
RecordDeserializationException是在Kafka客户端拉取线程中抛出的,未进入Spring Kafka的监听器执行链路,DefaultErrorHandler只能处理监听器执行阶段的异常。而ErrorHandlingDeserializer作为代理反序列化器,会捕获反序列化异常,将消息封装成包含异常的ConsumerRecord后传递给监听器,此时异常才能被DefaultErrorHandler捕获处理。
内容的提问来源于stack exchange,提问作者Higher-Kinded Type

