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

Spring Boot中控制Azure Service Bus消息监听器启停的实现方案

问题

已在Spring Boot应用中集成Azure Service Bus,当前应用启动时会立刻开始监听消息。现需要调整逻辑:默认禁用Azure Service Bus消息监听器,待ApplicationReadyEvent事件触发并完成指定任务后,再启用监听器接收队列或主题的消息。现有配置代码如下,需实现该需求:

application.yml

spring:
  cloud:
    azure:
      servicebus:
        namespace: **********
        
xxx:
  azure:
    servicebus:
      connection: ***********
      queue: **********

AzureConfiguration.java

import com.azure.spring.integration.servicebus.inbound.ServiceBusInboundChannelAdapter;
import com.azure.spring.messaging.servicebus.core.ServiceBusProcessorFactory;
import com.azure.spring.messaging.servicebus.core.listener.ServiceBusMessageListenerContainer;
import com.azure.spring.messaging.servicebus.core.properties.ServiceBusContainerProperties;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.integration.channel.DirectChannel;
import org.springframework.integration.dsl.IntegrationFlow;
import org.springframework.integration.dsl.IntegrationFlows;
import org.springframework.messaging.MessageChannel;


@Configuration
public class AzureConfiguration{

    @Value("${xxx.azure.servicebus.connection}")
    private String serviceBusConnection;

    @Value("${xxx.azure.servicebus.queue}")
    private String serviceBusQueue;

    private static final String SERVICE_BUS_INPUT_CHANNEL = "yyyyy";
    private static final String SENSOR_DATA_CHANNEL = "zzzzz";
    private static final String SERVICE_BUS_LISTENER_CONTAINER = "aaaaa";


    @Bean(name = SERVICE_BUS_LISTENER_CONTAINER)
    public ServiceBusMessageListenerContainer serviceBusMessageListenerContainer(ServiceBusProcessorFactory processorFactory) {

        ServiceBusContainerProperties containerProperties = new ServiceBusContainerProperties();
        containerProperties.setConnectionString(serviceBusConnection);
        containerProperties.setEntityName(serviceBusQueue);
        containerProperties.setAutoComplete(true);
        return new ServiceBusMessageListenerContainer(processorFactory, containerProperties);
    }


    @Bean
    public ServiceBusInboundChannelAdapter serviceBusInboundChannelAdapter(
            @Qualifier(SERVICE_BUS_INPUT_CHANNEL) MessageChannel inputChannel,
            @Qualifier(SERVICE_BUS_LISTENER_CONTAINER) ServiceBusMessageListenerContainer listenerContainer) {

        ServiceBusInboundChannelAdapter adapter = new ServiceBusInboundChannelAdapter(listenerContainer);
        adapter.setOutputChannel(inputChannel);
        
        return adapter;
    }


    @Bean(name = SERVICE_BUS_INPUT_CHANNEL)
    public MessageChannel serviceBusInputChannel() {

        return new DirectChannel();
    }


    @Bean(name = SENSOR_DATA_CHANNEL)
    public MessageChannel sensorDataChannel() {

        return new DirectChannel();
    }


    @Bean
    public IntegrationFlow serviceBusMessageFlow() {

        return IntegrationFlows.from(SERVICE_BUS_INPUT_CHANNEL)
                .<byte[], String>transform(String::new)
                .channel(SENSOR_DATA_CHANNEL)
                .get();
    }
}

AppEventListenerService.java

import com.azure.spring.integration.servicebus.inbound.ServiceBusInboundChannelAdapter;
import lombok.AllArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.context.event.ApplicationReadyEvent;
import org.springframework.context.event.EventListener;
import org.springframework.stereotype.Service;

import java.util.List;

@Slf4j
@Service
@AllArgsConstructor
public class AppEventListenerService{

   
    @EventListener(ApplicationReadyEvent.class)
    public void OnApplicationStarted() {
        log.debug("Enter OnApplicationStarted");
        // 默认禁用Azure Service Bus消息监听器
        // 执行一些任务
        // 启用Azure Service Bus消息监听器
        log.debug("Exit OnApplicationStarted");
    }
}
解决方案

要实现该需求,核心是让监听器容器默认不自动启动,待ApplicationReadyEvent触发并完成指定任务后再手动启动,具体修改如下:

1. 修改监听器容器配置,默认禁用自动启动

在AzureConfiguration的serviceBusMessageListenerContainer方法中,给容器添加setAutoStartup(false)配置,阻止其在应用启动阶段自动运行:

@Bean(name = SERVICE_BUS_LISTENER_CONTAINER)
public ServiceBusMessageListenerContainer serviceBusMessageListenerContainer(ServiceBusProcessorFactory processorFactory) {

    ServiceBusContainerProperties containerProperties = new ServiceBusContainerProperties();
    containerProperties.setConnectionString(serviceBusConnection);
    containerProperties.setEntityName(serviceBusQueue);
    containerProperties.setAutoComplete(true);
    
    ServiceBusMessageListenerContainer container = new ServiceBusMessageListenerContainer(processorFactory, containerProperties);
    // 默认禁用自动启动
    container.setAutoStartup(false);
    return container;
}

2. 在ApplicationReadyEvent事件中启用监听器

在AppEventListenerService中注入ServiceBusMessageListenerContainer,完成初始化任务后调用start()方法启动监听器:

import com.azure.spring.messaging.servicebus.core.listener.ServiceBusMessageListenerContainer;
import lombok.AllArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.context.event.ApplicationReadyEvent;
import org.springframework.context.event.EventListener;
import org.springframework.stereotype.Service;

@Slf4j
@Service
@AllArgsConstructor
public class AppEventListenerService{

    // 注入监听器容器
    private final ServiceBusMessageListenerContainer serviceBusMessageListenerContainer;
   
    @EventListener(ApplicationReadyEvent.class)
    public void onApplicationStarted() {
        log.debug("Enter onApplicationStarted");
        // 执行指定初始化任务(此处替换为你的业务代码)
        log.debug("完成初始化任务,准备启用Azure Service Bus消息监听器");
        
        // 启用Azure Service Bus消息监听器
        serviceBusMessageListenerContainer.start();
        log.debug("Azure Service Bus消息监听器已启用,开始接收消息");
        
        log.debug("Exit onApplicationStarted");
    }
}

原理说明

  • setAutoStartup(false)会让监听器容器跳过应用启动阶段的自动初始化流程,保持未激活状态;
  • ApplicationReadyEvent是Spring Boot在应用完全启动、所有Bean加载完成后触发的事件,此时执行初始化任务能确保依赖资源全部就绪;
  • 调用serviceBusMessageListenerContainer.start()会触发监听器容器启动,开始从Azure Service Bus队列/主题接收消息。

内容的提问来源于stack exchange,提问作者Kush Patel

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 22:40:56