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
相关产品推荐
相关产品推荐

