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

