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

如何在Spring Boot中读取ActiveMQ队列的未确认待处理消息

在Spring Boot中读取ActiveMQ队列的未确认(Pending)消息

我懂你的需求——你不想一有消息发进队列就立刻被消费,而是希望能主动读取队列里那些还没被确认的pending待处理消息,对吧?先说说你当前代码的问题:你用的@JmsListener是消息驱动的推模式,只要队列有新消息就会触发消费,这和你想要的「主动读取pending消息」逻辑完全不符。下面给你几种实用的实现方案:

方案一:用JmsTemplate主动拉取消息

这是最直接的方式——放弃推模式,改用拉模式主动去队列获取未被消费的消息,而且完全由你控制是否确认消息。

实现步骤:

  1. 确保你已经注入了JmsTemplate(Spring Boot默认会自动配置,没自定义的话直接用就行)
  2. 写一个主动拉取的方法:
@Autowired
private JmsTemplate jmsTemplate;

public void fetchPendingMessages() {
    // 从指定队列拉取一条pending消息,默认阻塞,也可以加超时时间(比如5000毫秒)
    Message message = jmsTemplate.receive("LOCAL.TEST");
    // Message message = jmsTemplate.receive("LOCAL.TEST", 5000); // 5秒超时返回null

    if (message != null) {
        try {
            // 解析消息内容,这里以TextMessage为例,可根据实际类型调整
            String messageContent = ((TextMessage) message).getText();
            System.out.println("拉取到pending消息:" + messageContent);
            
            // 注意:不要立刻调用acknowledge(),除非确认处理完成
            // message.acknowledge(); 
            
            // 不确认的话,这条消息会一直处于pending状态,重启消费端后仍能被拉取
        } catch (JMSException e) {
            e.printStackTrace();
        }
    }
}

方案二:自定义Listener容器,控制消费时机

如果你还是想保留@JmsListener的Listener模式,但希望能控制什么时候开始消费pending消息,可以通过自定义容器工厂+启停容器来实现。

实现步骤:

  1. 配置手动确认的JmsListener容器工厂:
@Bean
public JmsListenerContainerFactory<?> myJmsListenerContainerFactory(ConnectionFactory connectionFactory) {
    DefaultJmsListenerContainerFactory factory = new DefaultJmsListenerContainerFactory();
    factory.setConnectionFactory(connectionFactory);
    // 设置为客户端手动确认模式
    factory.setSessionAcknowledgeMode(Session.CLIENT_ACKNOWLEDGE);
    // 关闭消息预取,避免容器提前把消息拉到本地缓存(确保看到的是队列真实pending消息)
    factory.setPrefetchSize(0);
    // 关闭会话缓存,避免复用会话导致的状态问题
    factory.setCacheLevel(DefaultMessageListenerContainer.CACHE_NONE);
    return factory;
}
  1. 给Listener加上唯一ID,方便后续控制容器:
@JmsListener(destination = "LOCAL.TEST", containerFactory = "myJmsListenerContainerFactory", id = "testQueueListener")
public void receiveMessage(final Message jsonMessage) throws JMSException {
    LOGGER.info("=== 接收到pending消息 {}", jsonMessage);
    // 不要立即确认,处理完成后再调用下面的代码
    // jsonMessage.acknowledge();
}
  1. 注入容器注册表,控制Listener的启停:
@Autowired
private JmsListenerEndpointRegistry jmsListenerEndpointRegistry;

// 暂停Listener,不再接收新消息,队列消息保持pending状态
public void pauseQueueListener() {
    MessageListenerContainer container = jmsListenerEndpointRegistry.getListenerContainer("testQueueListener");
    if (container != null && container.isRunning()) {
        container.stop();
    }
}

// 重启Listener,开始拉取队列里的pending消息
public void resumeQueueListener() {
    MessageListenerContainer container = jmsListenerEndpointRegistry.getListenerContainer("testQueueListener");
    if (container != null && !container.isRunning()) {
        container.start();
    }
}

方案三:用ActiveMQ Admin API查询pending消息(仅查看不消费)

如果你只是想查看队列里有哪些pending消息,不需要消费它们,可以用ActiveMQ的QueueBrowser来实现——它只会遍历消息,不会改变消息的状态。

实现代码:

@Autowired
private ConnectionFactory connectionFactory;

public void queryPendingMessages() throws JMSException {
    // 强转为ActiveMQConnection,使用它的Admin API
    ActiveMQConnection connection = (ActiveMQConnection) connectionFactory.createConnection();
    connection.start();
    // 创建非事务性、自动确认的会话(这里的自动确认仅针对Browser,不影响消息本身状态)
    ActiveMQSession session = (ActiveMQSession) connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
    
    // 创建目标队列的浏览器
    Queue targetQueue = session.createQueue("LOCAL.TEST");
    QueueBrowser queueBrowser = session.createBrowser(targetQueue);
    
    // 遍历队列中的所有pending消息
    Enumeration<?> messageEnum = queueBrowser.getEnumeration();
    while (messageEnum.hasMoreElements()) {
        Message message = (Message) messageEnum.nextElement();
        System.out.println("Pending消息ID: " + message.getJMSMessageID());
        System.out.println("消息发送时间: " + message.getJMSTimestamp());
        // 可根据消息类型解析内容,比如TextMessage
        if (message instanceof TextMessage) {
            System.out.println("消息内容: " + ((TextMessage) message).getText());
        }
    }
    
    // 记得关闭资源
    queueBrowser.close();
    session.close();
    connection.close();
}

最后给你提个小建议:

  • 如果是要消费pending消息,优先用方案一(JmsTemplate拉取),逻辑更清晰;
  • 如果要保留Listener模式,用方案二;
  • 如果只是监控查询,用方案三就够了。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 03:34:32