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

如何在Message Listener中查看ActiveMQ队列大小?

在ActiveMQ Message Listener中获取队列大小的可行方案

问题背景

我在使用Message Listener处理ActiveMQ队列消息时,需要在监听过程中查看队列的当前大小。目前在Listener外部调用判断队列是否为空的逻辑可以正常工作,但在Listener内部无法实现这个需求,现有代码如下:

public class TestClass {
    MessageConsumer consumerTransformation;
    MessageListener listenerObjectTransformation;
    public static void main(String []args) throws JMSException {
        ActiveMQModel activeMQModelObject = new ActiveMQModel();
        //String subject = "TRANSFORMATION_QUEUE";
        String subject = "IMPORT_QUEUE";
         //consumerTransformation = activeMQModelObject.getActiveMQConsumer(subject);

          // Here we set the listener to listen to all the messages in the queue
          //listenerObjectTransformation = new TransformationMessageListener();
          //consumerTransformation.setMessageListener(listenerObjectTransformation);
          boolean isQueueEmpty = activeMQModelObject.isMessageQueueEmpty(subject);
          System.out.println("Size " + isQueueEmpty);
    }
    /*private class TransformationMessageListener implements MessageListener {
        @Override
        public void onMessage(Message messagearg) {
            System.out.println("test....");
        }
    }*/
}

可行解决方案

方法1:通过JMX获取队列统计信息

ActiveMQ默认开启JMX监控,可在Listener内部通过JMX连接到Broker,直接获取指定队列的实时大小:

private class TransformationMessageListener implements MessageListener {
    @Override
    public void onMessage(Message messagearg) {
        try {
            // 连接ActiveMQ的JMX服务(默认端口1099)
            JMXServiceURL jmxUrl = new JMXServiceURL("service:jmx:rmi:///jndi/rmi://localhost:1099/jmxrmi");
            JMXConnector jmxConnector = JMXConnectorFactory.connect(jmxUrl);
            MBeanServerConnection mBeanConn = jmxConnector.getMBeanServerConnection();
            
            // 构造队列的ObjectName,需替换为你的Broker名称和队列名
            ObjectName queueMBean = new ObjectName("org.apache.activemq:type=Broker,brokerName=localhost,destinationType=Queue,destinationName=IMPORT_QUEUE");
            // 获取队列当前消息数
            Long currentQueueSize = (Long) mBeanConn.getAttribute(queueMBean, "QueueSize");
            System.out.println("当前队列大小:" + currentQueueSize);
            
            jmxConnector.close();
        } catch (Exception e) {
            e.printStackTrace();
        }
        // 消息处理逻辑
        System.out.println("test....");
    }
}

注意:需确保Listener所在应用能访问Broker的JMX端口,若Broker配置了JMX认证,需添加对应的认证参数。

方法2:在Listener中注入队列操作实例

将ActiveMQModel实例通过构造函数传入Listener,直接复用外部可用的队列查询逻辑,同时保证线程安全:

public class TestClass {
    MessageConsumer consumerTransformation;
    MessageListener listenerObjectTransformation;
    private ActiveMQModel activeMQModelObject; // 改为类成员变量

    public static void main(String []args) throws JMSException {
        TestClass testInstance = new TestClass();
        testInstance.activeMQModelObject = new ActiveMQModel();
        String subject = "IMPORT_QUEUE";
        testInstance.consumerTransformation = testInstance.activeMQModelObject.getActiveMQConsumer(subject);
        // 初始化Listener时传入ActiveMQModel和队列名
        testInstance.listenerObjectTransformation = new TransformationMessageListener(testInstance.activeMQModelObject, subject);
        testInstance.consumerTransformation.setMessageListener(testInstance.listenerObjectTransformation);
    }

    private class TransformationMessageListener implements MessageListener {
        private final ActiveMQModel activeMQModel;
        private final String queueName;

        // 构造函数注入依赖
        public TransformationMessageListener(ActiveMQModel activeMQModel, String queueName) {
            this.activeMQModel = activeMQModel;
            this.queueName = queueName;
        }

        @Override
        public void onMessage(Message messagearg) {
            try {
                // 调用外部可用的队列空判断逻辑
                boolean isQueueEmpty = activeMQModel.isMessageQueueEmpty(queueName);
                // 若需要具体大小,可修改ActiveMQModel添加getQueueSize方法
                // Long queueSize = activeMQModel.getQueueSize(queueName);
                System.out.println("队列是否为空:" + isQueueEmpty);
            } catch (JMSException e) {
                e.printStackTrace();
            }
            // 消息处理逻辑
            System.out.println("test....");
        }
    }
}

说明:需确保ActiveMQModel中的队列查询方法是线程安全的,避免多线程调用时出现异常。

方法3:使用QueueBrowser查询队列大小

在Listener内部创建QueueBrowser,遍历队列统计消息数(注意:这种方式会遍历所有未消费消息,大队列下性能较低):

private class TransformationMessageListener implements MessageListener {
    private final Connection connection;
    private final String queueName;

    public TransformationMessageListener(Connection connection, String queueName) {
        this.connection = connection;
        this.queueName = queueName;
    }

    @Override
    public void onMessage(Message messagearg) {
        try (Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
             QueueBrowser browser = session.createBrowser(session.createQueue(queueName))) {
            
            Enumeration<?> messages = browser.getEnumeration();
            int count = 0;
            while (messages.hasMoreElements()) {
                messages.nextElement();
                count++;
            }
            System.out.println("当前队列大小:" + count);
            
        } catch (JMSException e) {
            e.printStackTrace();
        }
        System.out.println("test....");
    }
}

注意:该方法仅适合小队列场景,大队列下会导致性能问题。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 18:50:27