【问题标题】:ActiveMQ JMSListenerActiveMQ JMSListener
【发布时间】:2017-05-13 00:01:56
【问题描述】:

我是这个话题的新手。我正在使用 JMS 侦听器来侦听包含拆分消息的 Active MQ。我需要听队列直到最后一条消息,然后将它一起发送到 UI。我能够收听队列并获取消息,但我不知道有多少拆分消息可用,因此我无法将它们全部发送。有没有办法让监听器做上述操作?就像队列中没有可用的消息一样,jms 侦听器会产生空值吗?任何想法或帮助都会非常有帮助。

我正在使用下面的代码通过 JMS 监听器来监听队列。

 private static final String ORDER_RESPONSE_QUEUE = "mail-response-queue";

@JmsListener(destination = ORDER_RESPONSE_QUEUE)
public void receiveMessage(final Message<InventoryResponse> message) throws JMSException {
    LOG.info("+++++++++++++++++++++++++++++++++++++++++++++++++++++");
    MessageHeaders headers =  message.getHeaders();
    LOG.info("Application : headers received : {}", headers);

    InventoryResponse response = message.getPayload();
    LOG.info("Application : response received : {}",response);  
    LOG.info("+++++++++++++++++++++++++++++++++++++++++++++++++++++");
}

我可以使用 JMS 监听器获取队列信息吗?

【问题讨论】:

    标签: java jms activemq spring-integration spring-jms


    【解决方案1】:

    通过 jmx,您可以访问有关目的地的信息,例如,您可以了解消息在队列中的待处理情况。

    请注意,如果发送新消息,这可能会发生变化

    长 org.apache.activemq.broker.jmx.DestinationViewMBean.getQueueSize()

    @MBeanInfo(value="目标中尚未发送的消息数 被消耗。可能已发送但未确认。")

    返回此目标中尚未发送的消息数 消费返回:返回此目的地的消息数 还没吃完

    import java.util.HashMap;
    import java.util.Map;
    
    import javax.management.MBeanServerConnection;
    import javax.management.MBeanServerInvocationHandler;
    import javax.management.ObjectName;
    import javax.management.remote.JMXConnector;
    import javax.management.remote.JMXConnectorFactory;
    import javax.management.remote.JMXServiceURL;
    
    import org.apache.activemq.broker.jmx.BrokerViewMBean;
    import org.apache.activemq.broker.jmx.QueueViewMBean;
    
    public class JMXGetDestinationInfos {
    
        public static void main(String[] args) throws Exception {
            JMXServiceURL url = new JMXServiceURL("service:jmx:rmi:///jndi/rmi://host:1099/jmxrmi");
            Map<String, String[]> env = new HashMap<>();
            String[] creds = {"admin", "activemq"};
            env.put(JMXConnector.CREDENTIALS, creds);
            JMXConnector jmxc = JMXConnectorFactory.connect(url, env);
            MBeanServerConnection conn = jmxc.getMBeanServerConnection();
    
            ObjectName activeMq = new ObjectName("org.apache.activemq:type=Broker,brokerName=localhost");
    
            BrokerViewMBean mbean = MBeanServerInvocationHandler.newProxyInstance(conn, activeMq, BrokerViewMBean.class,
                    true);
            for (ObjectName name : mbean.getQueues()) {
                if (("Destination".equals(name.getKeyProperty("destinationName")))) {
                    QueueViewMBean queueMbean = MBeanServerInvocationHandler.newProxyInstance(conn, name,
                            QueueViewMBean.class, true);
                    System.out.println(queueMbean.getQueueSize());
                }
            }
        }
    }
    

    为什么不消费消息,当没有收到消息时显示?如果没有收到消息,您可以使用下面的方法在超时后返回 null。

    ActiveMQMessageConsumer.receive(长时间超时) throws JMSException 接收在指定超时间隔内到达的下一条消息。此调用阻塞,直到 一条消息到达,超时到期,或者这个消息消费者是 关闭。零超时永不过期,调用阻塞 无限期地。指定者:接口MessageConsumer中的receive 参数: timeout - 超时值(以毫秒为单位),超时 零永不过期。返回:为此生成的下一条消息 消息使用者,如果超时过期或此消息,则返回 null 消费者同时关闭

    更新

    可能是这样的:

    import java.io.IOException;
    import java.net.MalformedURLException;
    import java.util.HashMap;
    import java.util.Map;
    
    import javax.management.MBeanServerConnection;
    import javax.management.MBeanServerInvocationHandler;
    import javax.management.MalformedObjectNameException;
    import javax.management.ObjectName;
    import javax.management.remote.JMXConnector;
    import javax.management.remote.JMXConnectorFactory;
    import javax.management.remote.JMXServiceURL;
    
    import org.apache.activemq.broker.jmx.BrokerViewMBean;
    import org.apache.activemq.broker.jmx.QueueViewMBean;
    
    public class JMXGetDestinationInfos {
    
        private QueueViewMBean queueMbean;
    
        {
            try {
                JMXServiceURL url = new JMXServiceURL("service:jmx:rmi:///jndi/rmi://host:1099/jmxrmi");
                Map<String, String[]> env = new HashMap<>();
                String[] creds = { "admin", "activemq" };
                env.put(JMXConnector.CREDENTIALS, creds);
                JMXConnector jmxc = JMXConnectorFactory.connect(url, env);
                MBeanServerConnection conn = jmxc.getMBeanServerConnection();
    
                ObjectName activeMq = new ObjectName("org.apache.activemq:type=Broker,brokerName=localhost");
    
                BrokerViewMBean mbean = MBeanServerInvocationHandler.newProxyInstance(conn, activeMq, BrokerViewMBean.class,
                        true);
                for (ObjectName name : mbean.getQueues()) {
                    if (("Destination".equals(name.getKeyProperty("destinationName")))) {
                        queueMbean = MBeanServerInvocationHandler.newProxyInstance(conn, name, QueueViewMBean.class, true);
                        System.out.println(queueMbean.getQueueSize());
                        break;
                    }
                }
            } catch (Exception e) {
                e.printStackTrace();
            }
        }
    
        @JmsListener(destination = ORDER_RESPONSE_QUEUE)
        public void receiveMessage(final Message<InventoryResponse> message, javax.jms.Message amqMessage) throws JMSException {
            LOG.info("+++++++++++++++++++++++++++++++++++++++++++++++++++++");
            MessageHeaders headers =  message.getHeaders();
            LOG.info("Application : headers received : {}", headers);
    
            InventoryResponse response = message.getPayload();
            LOG.info("Application : response received : {}",response);  
            LOG.info("+++++++++++++++++++++++++++++++++++++++++++++++++++++");
            //queueMbean.getQueueSize()  is real time, each call return the real size
            ((org.apache.activemq.command.ActiveMQMessage) amqMessage ).acknowledge();
            if(queueMbean != null && queueMbean.getQueueSize() == 0){
                //display messages ??
            }
        }
    }
    

    因为getQueueSize() 返回的消息数量 尚未消费的目的地。 可能已派出但 未确认。

    一种解决方案是在春季 DefaultMessageListenerContainer.sessionAcknowledgeModeName 中将确认模式更新为 org.apache.activemq.ActiveMQSession.INDIVIDUAL_ACKNOWLEDGE 以创建会话,并单独确认每条消息,然后检查大小是否 == 0(大小 == 0 表示所有消息都已发送并确认)。

    【讨论】:

    • 谢谢,但我想使用监听器而不是消费者,是否可以使用 JMS 监听器获取队列信息?
    • 我已经更新了我的帖子。请看一下并帮助我
    猜你喜欢
    • 2021-08-09
    • 2019-07-30
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-01-04
    • 2014-11-12
    相关资源
    最近更新 更多