Spring 2.x迁移3.x后使用DefaultErrorHandler遇KafkaException处理异常
Spring Kafka 3.x迁移:DefaultErrorHandler报错及替代SeekToCurrentBatchErrorHandler方案
背景
- 项目从Spring 2.x迁移至3.x,同步升级
spring-kafka到3.0.13版本 - 原使用的
SeekToCurrentBatchErrorHandler已被弃用,替换为以下配置:factory.setCommonErrorHandler(new DefaultErrorHandler(new FixedBackOff(2000, 5)));
报错信息
java.lang.IllegalStateException: This error handler cannot process 'org.apache.kafka.common.KafkaException's; no record information is available at org.springframework.kafka.listener.DefaultErrorHandler.handleOtherException(DefaultErrorHandler.java:204) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.handleConsumerException(KafkaMessageListenerContainer.java:1961) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:1396) at java.base/java.util.concurrent.CompletableFuture$AsyncRun.run(CompletableFuture.java:1804) at java.base/java.lang.Thread.run(Thread.java:840) Caused by: org.apache.kafka.common.KafkaException: Unexpected error from SyncGroup: The server experienced an unexpected error when processing the request. at org.apache.kafka.clients.consumer.internals.AbstractCoordinator$SyncGroupResponseHandler.handle(AbstractCoordinator.java:852) at org.apache.kafka.clients.consumer.internals.AbstractCoordinator$SyncGroupResponseHandler.handle(AbstractCoordinator.java:771) at org.apache.kafka.clients.consumer.internals.AbstractCoordinator$CoordinatorResponseHandler.onSuccess(AbstractCoordinator.java:1260) at org.apache.kafka.clients.consumer.internals.AbstractCoordinator$CoordinatorResponseHandler.onSuccess(AbstractCoordinator.java:1235) at org.apache.kafka.clients.consumer.internals.RequestFuture$1.onSuccess(RequestFuture.java:206) at org.apache.kafka.clients.consumer.internals.RequestFuture.fireSuccess(RequestFuture.java:169) at org.apache.kafka.clients.consumer.internals.RequestFuture.complete(RequestFuture.java:129) at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient$RequestFutureCompletionHandler.fireCompletion(ConsumerNetworkClient.java:617) at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.firePendingCompletedRequests(ConsumerNetworkClient.java:427) at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.poll(ConsumerNetworkClient.java:312) at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.poll(ConsumerNetworkClient.java:251) at org.apache.kafka.clients.consumer.KafkaConsumer.pollForFetches(KafkaConsumer.java:1307) at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1243) at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1216) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollConsumer(KafkaMessageListenerContainer.java:1676) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doPoll(KafkaMessageListenerContainer.java:1651) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollAndInvoke(KafkaMessageListenerContainer.java:1452) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:1344) ... 2 more
依赖配置
<dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId> <version>3.0.13</version> </dependency>
需求
新错误处理器需具备与SeekToCurrentBatchErrorHandler完全一致的业务逻辑,解决上述报错问题。
解决方案
报错原因
DefaultErrorHandler默认无法处理无记录关联的Kafka协调层异常(如SyncGroup错误),而原SeekToCurrentBatchErrorHandler会自动重启消费者容器来应对这类异常。同时,DefaultErrorHandler默认按单条记录处理,若使用批量监听器,需额外配置开启批量支持。
配置实现(与原处理器逻辑对齐)
// 定义回退策略:2秒间隔,最多重试5次 FixedBackOff fixedBackOff = new FixedBackOff(2000, 5); // 自定义DefaultErrorHandler,对齐SeekToCurrentBatchErrorHandler逻辑 DefaultErrorHandler errorHandler = new DefaultErrorHandler(fixedBackOff) { @Override public void handleOtherException(Exception thrownException, Consumer<?, ?> consumer, MessageListenerContainer container, boolean batchListener) { // 模拟原处理器行为:重启消费者容器以恢复协调层异常 container.stop(); container.start(); } }; // 关闭处理后自动提交偏移量,与原批处理器逻辑一致(重试失败后才提交) errorHandler.setAckAfterHandle(false); // 若使用批量监听器,必须开启批量模式支持 errorHandler.setBatchListener(true); // 配置到消费者工厂 factory.setCommonErrorHandler(errorHandler); // 批量监听器需额外设置(若使用) factory.setBatchListener(true);
关键说明
- handleOtherException重写:处理无记录的Kafka异常,通过重启容器恢复,和原
SeekToCurrentBatchErrorHandler行为一致。 - setAckAfterHandle(false):原批处理器仅在重试失败后才提交偏移量,关闭自动提交可对齐该逻辑。
- 批量监听器配置:若你的
@KafkaListener接收List<ConsumerRecord>类型参数,必须设置factory.setBatchListener(true),否则DefaultErrorHandler会按单条记录处理,偏离原逻辑。
内容的提问来源于stack exchange,提问作者Chandan Bansal
相关产品推荐
相关产品推荐

