Spring Boot中无法动态停止IBM MQ JMS消费者问题求助
无法通过Spring Boot REST端点动态停止IBM MQ JMS消费者
环境信息
IBM MQ版本:9.2.0.5
相关代码
pom.xml
<dependency> <groupId>com.ibm.mq</groupId> <artifactId>mq-jms-spring-boot-starter</artifactId> <version>2.0.8</version> </dependency>
JmsConfig.java
@Configuration @EnableJms @Log4j2 public class JmsConfig { @Bean public MQQueueConnectionFactory mqQueueConnectionFactory() { MQQueueConnectionFactory mqQueueConnectionFactory = new MQQueueConnectionFactory(); mqQueueConnectionFactory.setHostName("my-ibm-mq-host.com"); try { mqQueueConnectionFactory.setTransportType(WMQConstants.WMQ_CM_CLIENT); mqQueueConnectionFactory.setCCSID(1208); mqQueueConnectionFactory.setChannel("my-channel"); mqQueueConnectionFactory.setPort(1234); mqQueueConnectionFactory.setQueueManager("my-QM"); } catch (Exception e) { log.error("Exception while creating JMS connecion...", e.getMessage()); } return mqQueueConnectionFactory; } }
JmsListenerConfig.java
@Configuration @Log4j2 public class JmsListenerConfig implements JmsListenerConfigurer { @Autowired private JmsConfig jmsConfig; private Map<String, String> queueMap = new HashMap<>(); @Bean public DefaultJmsListenerContainerFactory mqJmsListenerContainerFactory() throws JMSException { DefaultJmsListenerContainerFactory factory = new DefaultJmsListenerContainerFactory(); factory.setConnectionFactory(jmsConfig.mqQueueConnectionFactory()); factory.setDestinationResolver(new DynamicDestinationResolver()); factory.setSessionTransacted(true); factory.setConcurrency("5"); return factory; } @Override public void configureJmsListeners(JmsListenerEndpointRegistrar registrar) { queueMap.put("my-queue-101", "101"); log.info("queueMap: " + queueMap); queueMap.entrySet().forEach(e -> { SimpleJmsListenerEndpoint endpoint = new SimpleJmsListenerEndpoint(); endpoint.setDestination(e.getKey()); endpoint.setId(e.getValue()); try { log.info("Reading message...."); endpoint.setMessageListener(message -> { try { log.info("Receieved ID: {} Destination {}", message.getJMSMessageID(), message.getJMSDestination()); } catch (JMSException ex) { log.error("Exception while reading message - " + ex.getMessage()); } }); registrar.setContainerFactory(mqJmsListenerContainerFactory()); } catch (JMSException ex) { log.error("Exception while reading message - " + ex.getMessage()); } registrar.registerEndpoint(endpoint); }); } }
JmsController.java
@RestController @RequestMapping("/jms") @Log4j2 public class JmsController { @Autowired ApplicationContext context; @RequestMapping(value = "/stop", method = RequestMethod.GET) public @ResponseBody String haltJmsListener() { JmsListenerEndpointRegistry listenerEndpointRegistry = context.getBean(JmsListenerEndpointRegistry.class); Set<String> containerIds = listenerEndpointRegistry.getListenerContainerIds(); log.info("containerIds: " + containerIds); //stops all consumers listenerEndpointRegistry.stop(); //DOESN'T WORK :( //stops a consumer by id, used when there are multiple consumers and want to stop them individually //listenerEndpointRegistry.getListenerContainer("101").stop(); //DOESN'T WORK EITHER :( return "Jms Listener stopped"; } }
观察现象
- 初始消费者数量:0(符合预期)
- 服务器启动并建立队列连接后,消费者总数:1(符合预期)
- 调用
http://localhost:8080/jms/stop端点后,消费者总数:1(不符合预期,应回到0)
问题排查与修复方案
1. 修复容器工厂重复创建问题
在configureJmsListeners方法中,每次循环注册端点时都调用mqJmsListenerContainerFactory()创建新的工厂实例,导致每个端点绑定的容器工厂不是同一个Spring管理的Bean,注册表无法正确追踪和管理这些容器。
修改后的JmsListenerConfig.java:
@Configuration @Log4j2 public class JmsListenerConfig implements JmsListenerConfigurer { @Autowired private JmsConfig jmsConfig; @Autowired private DefaultJmsListenerContainerFactory mqJmsListenerContainerFactory; private Map<String, String> queueMap = new HashMap<>(); @Bean public DefaultJmsListenerContainerFactory mqJmsListenerContainerFactory() throws JMSException { DefaultJmsListenerContainerFactory factory = new DefaultJmsListenerContainerFactory(); factory.setConnectionFactory(jmsConfig.mqQueueConnectionFactory()); factory.setDestinationResolver(new DynamicDestinationResolver()); factory.setSessionTransacted(true); factory.setConcurrency("5"); return factory; } @Override public void configureJmsListeners(JmsListenerEndpointRegistrar registrar) { queueMap.put("my-queue-101", "101"); log.info("queueMap: " + queueMap); // 全局设置一次容器工厂,无需循环重复设置 registrar.setContainerFactory(mqJmsListenerContainerFactory); queueMap.entrySet().forEach(e -> { SimpleJmsListenerEndpoint endpoint = new SimpleJmsListenerEndpoint(); endpoint.setDestination(e.getKey()); endpoint.setId(e.getValue()); try { log.info("Reading message...."); endpoint.setMessageListener(message -> { try { log.info("Receieved ID: {} Destination {}", message.getJMSMessageID(), message.getJMSDestination()); } catch (JMSException ex) { log.error("Exception while reading message - " + ex.getMessage()); } }); } catch (Exception ex) { log.error("Exception while reading message - " + ex.getMessage()); } registrar.registerEndpoint(endpoint); }); } }
2. 处理容器停止的异步特性
stop()方法是异步执行的,不会立即终止所有消费者。可以通过回调或者状态检查确认停止操作完成:
修改后的JmsController.java:
@RestController @RequestMapping("/jms") @Log4j2 public class JmsController { @Autowired ApplicationContext context; @RequestMapping(value = "/stop", method = RequestMethod.GET) public @ResponseBody String haltJmsListener() { JmsListenerEndpointRegistry listenerEndpointRegistry = context.getBean(JmsListenerEndpointRegistry.class); Set<String> containerIds = listenerEndpointRegistry.getListenerContainerIds(); log.info("containerIds: " + containerIds); // 异步停止并通过回调确认完成 listenerEndpointRegistry.stop(() -> log.info("所有JMS容器已完成停止")); // 验证单个容器状态 MessageListenerContainer container = listenerEndpointRegistry.getListenerContainer("101"); if (container != null) { container.stop(); log.info("容器101运行状态: " + container.isRunning()); } return "Jms Listener停止指令已发送"; } }
3. 考虑事务会话的影响
由于配置了setSessionTransacted(true),容器会等待当前事务提交后才会停止。如果有未处理完成的消息,容器会等待事务结束再终止,可以调整事务超时时间或者在停止前检查消息处理状态。
内容的提问来源于stack exchange,提问作者Nital
相关产品推荐
相关产品推荐

