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

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);

关键说明

  1. handleOtherException重写:处理无记录的Kafka异常,通过重启容器恢复,和原SeekToCurrentBatchErrorHandler行为一致。
  2. setAckAfterHandle(false):原批处理器仅在重试失败后才提交偏移量,关闭自动提交可对齐该逻辑。
  3. 批量监听器配置:若你的@KafkaListener接收List<ConsumerRecord>类型参数,必须设置factory.setBatchListener(true),否则DefaultErrorHandler会按单条记录处理,偏离原逻辑。

内容的提问来源于stack exchange,提问作者Chandan Bansal

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 08:25:54