Spring JMS:如何在ExceptionListener类中获取currentAttempts的值?
解决方案
方案1:自定义可跟踪重试次数的BackOff
Spring JMS的DefaultMessageListenerContainer依赖BackOff接口控制重试间隔,默认使用FixedBackOff。我们可以包装该接口实现,在每次重试时记录次数,让自定义ExceptionListener能直接获取这个计数:
- 实现带计数功能的
TrackableBackOff:
public class TrackableBackOff implements BackOff { private final BackOff delegate; private final AtomicInteger currentAttempts = new AtomicInteger(0); public TrackableBackOff(BackOff delegate) { this.delegate = delegate; } @Override public long nextBackOff() { long backOff = delegate.nextBackOff(); if (backOff != BackOff.STOP) { currentAttempts.incrementAndGet(); } else { currentAttempts.set(0); } return backOff; } public int getCurrentAttempts() { return currentAttempts.get(); } public void reset() { currentAttempts.set(0); } }
- 配置容器时注入自定义BackOff,并关联到ExceptionListener:
@Configuration public class JmsConfig { @Bean public TrackableBackOff trackableBackOff() { FixedBackOff fixedBackOff = new FixedBackOff(); fixedBackOff.setInterval(5000); fixedBackOff.setMaxAttempts(1); // 保持原有配置 return new TrackableBackOff(fixedBackOff); } @Bean public DefaultMessageListenerContainer jmsListenerContainer(ConnectionFactory connectionFactory, TrackableBackOff trackableBackOff, CustomExceptionListener exceptionListener) { DefaultMessageListenerContainer container = new DefaultMessageListenerContainer(); container.setConnectionFactory(connectionFactory); container.setDestinationName("mydestination"); container.setBackOff(trackableBackOff); container.setExceptionListener(exceptionListener); // 其他必要配置 return container; } @Bean public CustomExceptionListener customExceptionListener(TrackableBackOff trackableBackOff) { return new CustomExceptionListener(trackableBackOff); } }
- 在自定义ExceptionListener中使用计数触发告警:
public class CustomExceptionListener implements ExceptionListener { private final TrackableBackOff trackableBackOff; private static final int ALERT_THRESHOLD = 50; private boolean alertSent = false; // 防抖,避免重复发邮件 public CustomExceptionListener(TrackableBackOff trackableBackOff) { this.trackableBackOff = trackableBackOff; } @Override public void onException(JMSException exception) { int currentAttempts = trackableBackOff.getCurrentAttempts(); if (currentAttempts >= ALERT_THRESHOLD && !alertSent) { sendAlertEmail(currentAttempts); alertSent = true; } // 连接恢复时重置告警标记 if (currentAttempts == 0) { alertSent = false; } } private void sendAlertEmail(int attempts) { // 邮件发送逻辑实现 } }
方案2:扩展DefaultMessageListenerContainer传递重试次数
如果不想修改BackOff逻辑,可以直接扩展容器类,在连接刷新失败时主动将重试次数传递给ExceptionListener:
- 扩展容器并重写连接刷新方法:
public class TrackableMessageListenerContainer extends DefaultMessageListenerContainer { @Override protected void refreshConnectionUntilSuccessful() throws Throwable { BackOffExecution backOffExecution = getBackOff().start(); long backOff; int attemptCount = 0; do { attemptCount++; try { refreshConnection(); return; } catch (Exception ex) { publishConnectionFailedEvent(ex); // 将重试次数传递给自定义监听 if (getExceptionListener() instanceof TrackableExceptionListener) { ((TrackableExceptionListener) getExceptionListener()).onException(ex, attemptCount); } else if (getExceptionListener() != null) { getExceptionListener().onException(ex instanceof JMSException ? (JMSException) ex : new JMSException(ex.getMessage())); } } backOff = backOffExecution.nextBackOff(); if (backOff != BackOff.STOP) { Thread.sleep(backOff); } } while (backOff != BackOff.STOP); throw new IllegalStateException("连接重试" + attemptCount + "次后仍失败"); } }
- 定义带重试次数参数的监听接口:
public interface TrackableExceptionListener extends ExceptionListener { void onException(Exception ex, int currentAttempts); }
- 实现监听接口处理告警:
public class CustomTrackableExceptionListener implements TrackableExceptionListener { private static final int ALERT_THRESHOLD = 50; private boolean alertSent = false; @Override public void onException(JMSException exception) { // 保留原有接口默认实现 } @Override public void onException(Exception ex, int currentAttempts) { if (currentAttempts >= ALERT_THRESHOLD && !alertSent) { sendAlertEmail(currentAttempts); alertSent = true; } // 连接恢复时重置标记 if (currentAttempts == 0) { alertSent = false; } } private void sendAlertEmail(int attempts) { // 邮件发送逻辑实现 } }
- 配置时使用自定义容器:
@Bean public TrackableMessageListenerContainer jmsListenerContainer(ConnectionFactory connectionFactory, CustomTrackableExceptionListener exceptionListener) { TrackableMessageListenerContainer container = new TrackableMessageListenerContainer(); container.setConnectionFactory(connectionFactory); container.setDestinationName("mydestination"); container.setExceptionListener(exceptionListener); // 其他必要配置 return container; }
关键注意点
- 必须在broker恢复连接后重置计数和告警标记,避免下次异常时计数不准确。
- 加入防抖逻辑,防止达到阈值后重复发送告警邮件。
内容的提问来源于stack exchange,提问作者Developer
相关产品推荐
相关产品推荐

