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

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

我的疑问

  1. 为什么日志显示连接已恢复,但实际应用还是无法消费消息,必须重启才能解决?
  2. 我的配置或代码中有没有遗漏的点,导致自动恢复机制失效?
  3. 有没有什么配置调整或者代码优化,可以让连接断开后自动恢复正常消费,不用手动重启应用?

麻烦各位大佬帮忙看看,万分感谢!

内容来源于stack exchange

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.08 14:44:31