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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 07:31:17