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

Spring Kafka监听遇JdbcException:重试后停止消费的实现方案

实现Spring Kafka重试后停止容器的错误处理逻辑

要实现「针对数据库超时(JdbcException)先重试,重试耗尽后停止Kafka监听器容器」的需求,无需直接组合DefaultErrorHandler和CommonContainerStoppingErrorHandler,只需基于DefaultErrorHandler扩展,在重试耗尽后添加停止容器的逻辑即可,具体实现如下:

方案1:自定义RecoveryCallback实现重试后停止容器

通过给DefaultErrorHandler配置自定义的RecoveryCallback,在重试次数用完后触发容器停止操作:

步骤1:配置自定义ErrorHandler

import org.springframework.kafka.listener.DefaultErrorHandler;
import org.springframework.kafka.listener.KafkaListenerEndpointRegistry;
import org.springframework.kafka.listener.MessageListenerContainer;
import org.springframework.util.backoff.FixedBackOff;

@Bean
public CommonErrorHandler jdbcRetryThenStopErrorHandler(KafkaListenerEndpointRegistry endpointRegistry) {
    // 配置重试策略:间隔1秒,重试3次
    FixedBackOff retryBackOff = new FixedBackOff(1000L, 3L);

    DefaultErrorHandler errorHandler = new DefaultErrorHandler(
        // 重试耗尽后的恢复逻辑
        (record, exception) -> {
            // 根据监听器ID获取对应容器并停止
            MessageListenerContainer container = endpointRegistry.getListenerContainer("your-db-write-listener");
            if (container != null && container.isRunning()) {
                container.stop();
                log.error("数据库超时重试耗尽,停止Kafka监听器容器,消息偏移量:{}", record.offset(), exception);
            }
        },
        retryBackOff
    );

    // 仅对JdbcException触发重试
    errorHandler.addRetryableExceptions(JdbcException.class);
    // 其他异常直接走恢复逻辑(不重试)
    errorHandler.addNotRetryableExceptions(Exception.class);

    return errorHandler;
}

步骤2:关联ErrorHandler到监听器

可以通过容器工厂全局配置,或者在单个@KafkaListener注解中指定:

方式A:全局容器工厂配置

@Bean
public ConcurrentKafkaListenerContainerFactory<String, YourMessageDto> kafkaListenerContainerFactory(
        ConsumerFactory<String, YourMessageDto> consumerFactory,
        CommonErrorHandler jdbcRetryThenStopErrorHandler) {
    ConcurrentKafkaListenerContainerFactory<String, YourMessageDto> factory =
            new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory);
    factory.setCommonErrorHandler(jdbcRetryThenStopErrorHandler);
    return factory;
}

方式B:单个监听器指定

@KafkaListener(
    topics = "your-topic",
    groupId = "db-write-group",
    id = "your-db-write-listener", // 需和ErrorHandler中指定的listenerId一致
    errorHandler = "jdbcRetryThenStopErrorHandler"
)
public void handleMessage(ConsumerRecord<String, YourMessageDto> record) {
    // 数据库写入业务逻辑
}

方案2:继承DefaultErrorHandler重写处理逻辑

如果需要更灵活的控制,可以直接继承DefaultErrorHandler,重写handleRemaining方法,在父类处理完成后添加容器停止逻辑:

步骤1:自定义ErrorHandler类

import org.springframework.kafka.listener.DefaultErrorHandler;
import org.springframework.kafka.listener.KafkaListenerEndpointRegistry;
import org.springframework.kafka.listener.MessageListenerContainer;
import org.springframework.util.backoff.FixedBackOff;

public class RetryThenStopErrorHandler extends DefaultErrorHandler {

    private final KafkaListenerEndpointRegistry endpointRegistry;
    private final String targetListenerId;

    public RetryThenStopErrorHandler(KafkaListenerEndpointRegistry endpointRegistry, String targetListenerId) {
        super(new FixedBackOff(1000L, 3L));
        this.endpointRegistry = endpointRegistry;
        this.targetListenerId = targetListenerId;
        // 仅对JdbcException重试
        addRetryableExceptions(JdbcException.class);
        addNotRetryableExceptions(Exception.class);
    }

    @Override
    protected void handleRemaining(Exception thrownException, ConsumerRecord<?, ?> record,
                                   Consumer<?, ?> consumer, MessageListenerContainer container) {
        // 先执行父类的默认处理(如记录日志)
        super.handleRemaining(thrownException, record, consumer, container);
        // 停止目标监听器容器
        MessageListenerContainer targetContainer = endpointRegistry.getListenerContainer(targetListenerId);
        if (targetContainer != null && targetContainer.isRunning()) {
            targetContainer.stop();
            log.error("JdbcException重试耗尽,停止监听器容器:{}", targetListenerId);
        }
    }
}

步骤2:配置自定义ErrorHandler

@Bean
public CommonErrorHandler retryThenStopErrorHandler(KafkaListenerEndpointRegistry endpointRegistry) {
    return new RetryThenStopErrorHandler(endpointRegistry, "your-db-write-listener");
}

关键注意事项

  • 监听器ID匹配:endpointRegistry.getListenerContainer中的ID必须和@KafkaListener注解的id属性完全一致,否则无法定位到目标容器。
  • 重试策略调整:可以根据业务需求替换FixedBackOff为ExponentialBackOff(指数退避),调整重试间隔和次数。
  • 多实例场景:停止容器仅会终止当前实例的消息消费,若需全局停止消费,需额外引入分布式协调机制(如配置中心开关)。
  • 容器恢复:容器停止后需手动重启,或通过监控平台触发自动恢复逻辑。

内容的提问来源于stack exchange,提问作者G Quintana

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 15:05:31