Spring Boot 3.4.6应用JMS Listener与Azure Service Bus随机断连且需重启恢复的问题求助
Spring Boot 3.4.6应用JMS Listener与Azure Service Bus随机断连且需重启恢复的问题求助
大家好,我最近遇到一个非常头疼的问题,想请各位帮忙分析排查下:
我维护的一个基于Spring Boot 3.4.6的应用,通过@JmsListener注解的方法接收Azure Service Bus主题的消息,但最近发现应用和Service Bus的连接会随机断开,而且必须重启应用才能重新恢复消息消费,靠组件自身的自动恢复机制完全没用。
断连时的关键错误日志
下面是连接断开时的核心日志片段:
2025-06-24T07:05:14.373Z INFO 1 --- [-f602ff55bc9e:1] org.apache.qpid.jms.JmsSession : A JMS MessageConsumer has been closed: JmsConsumerInfo: { ID:d85e265a-2c42-4c82-9ad1-f602ff55bc9e:1:9:2, destination = ***-topic } 2025-06-24T07:05:14.373Z WARN 1 --- [tContainer#0-14] o.s.j.l.DefaultMessageListenerContainer : Setup of JMS message listener invoker failed for destination '***-topic' - trying to recover. Cause: Unknown error from remote peer 2025-06-24T07:05:14.374Z INFO 1 --- [-f602ff55bc9e:1] org.apache.qpid.jms.JmsSession : A JMS MessageConsumer has been closed: JmsConsumerInfo: { ID:d85e265a-2c42-4c82-9ad1-f602ff55bc9e:1:10:2, destination = ***-topic } 2025-06-24T07:05:14.374Z WARN 1 --- [tContainer#1-14] o.s.j.l.DefaultMessageListenerContainer : Setup of JMS message listener invoker failed for destination '***-topic' - trying to recover. Cause: Unknown error from remote peer 2025-06-24T07:05:14.381Z INFO 1 --- [ync work thread] org.apache.qpid.jms.JmsConnection : Connection ID:d85e265a-2c42-4c82-9ad1-f602ff55bc9e:1 interrupted to server: amqps://***-prod.servicebus.windows.net 2025-06-24T07:05:14.451Z INFO 1 --- [ync work thread] org.apache.qpid.jms.JmsConnection : Connection ID:d85e265a-2c42-4c82-9ad1-f602ff55bc9e:1 restored to server: amqps://***-prod.servicebus.windows.net
从日志来看,虽然最后打印了连接已恢复的信息,但实际应用完全收不到任何消息,必须手动重启应用才能恢复正常消费。
我的相关配置与代码
1. application.yml 核心配置
spring: threads: virtual: enabled: true jms: listener: receive-timeout: 20s servicebus: listener: subscription-durable: true subscription-shared: true enabled: true connection-string: <connection-string> idle-timeout: 60000 pricing-tier: standard pool: enabled: true time-between-expiration-check: 5m
2. JMS Listener容器工厂Bean
@Bean(name = "customTopicJmsListenerContainerFactory") public @NonNull DefaultJmsListenerContainerFactory customTopicJmsListenerContainerFactory( @NonNull ConnectionFactory connectionFactory, @NonNull EventMessageConverter eventMessageConverter, @NonNull JmsErrorHandler jmsErrorHandler ) { var factory = new DefaultJmsListenerContainerFactory(); factory.setConnectionFactory(connectionFactory); factory.setMessageConverter(eventMessageConverter); factory.setErrorHandler(jmsErrorHandler); // 配置订阅共享,让服务所有实例都能连接到同一个订阅 factory.setSubscriptionShared(true); return factory; }
3. JMS错误处理器
public class JmsErrorHandler implements ErrorHandler { @Override public void handleError(@NonNull Throwable t) { log.warn("Error occurred on jms.", t); } }
4. 消息转换器
public class EventMessageConverter implements MessageConverter { private final ObjectMapper objectMapper; @Override public @NonNull Message toMessage( @NonNull Object object, @NonNull Session session ) { if (!(object instanceof Event<?> event)) { throw new IllegalArgumentException(String.format("Object must be an event class, provided: '%s'.", object)); } try { final var json = objectMapper.writeValueAsString(event); final var message = (JmsTextMessage) session.createTextMessage(json); message.setJMSMessageID(event.getId()); Optional.ofNullable(event.getScheduledAt()) .ifPresent(scheduledAt -> message.getFacade().setTracingAnnotation("x-opt-scheduled-enqueue-time", Date.from(scheduledAt))); return message; } catch (Exception e) { throw new UnsupportedOperationException(e); } } @Override public @NonNull Event<?> fromMessage(@NonNull Message message) { if (!(message instanceof TextMessage textMessage)) { throw new IllegalArgumentException(String.format("Message must be a text message, provided: '%s'.", message)); } try { return objectMapper.readValue(textMessage.getText(), Event.class); } catch (Exception e) { log.error("Cant read event message '{}'.", textMessage, e); throw new UnsupportedOperationException(e); } } }
5. 消息监听器类
public class AzureServiceBusListener implements EventListener { private final ApplicationEventPublisher publisher; private final List<EventFilter> eventFilters; @JmsListener( destination = "${<destination>}", subscription = "${<subscription>}", containerFactory = "customTopicJmsListenerContainerFactory" ) public void onEvent(@NonNull Event<?> event) { try { if (event instanceof IgnoredEvent) { log.info("Ignored event '{}' received.", event); return; } final var accepted = eventFilters.stream() .anyMatch(eventFilter -> eventFilter.accept(event)); if (!accepted) { log.info("Event ({}) received. Event is not accepted. Ignore it.", event.toReadableString()); return; } log.info("Event ({}) received. Propagate the event.", event.toReadableString()); publisher.publishEvent(event); } catch (Exception e) { log.error("Exception in AzureServiceBusListener when processing event", e); } } }
我的疑问
- 为什么日志显示连接已恢复,但实际应用还是无法消费消息,必须重启才能解决?
- 我的配置或代码中有没有遗漏的点,导致自动恢复机制失效?
- 有没有什么配置调整或者代码优化,可以让连接断开后自动恢复正常消费,不用手动重启应用?
麻烦各位大佬帮忙看看,万分感谢!
内容来源于stack exchange
相关产品推荐
相关产品推荐

