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

