【问题标题】:Does ActiveMQ support max message processing time in consumerActiveMQ 是否支持消费者的最大消息处理时间
【发布时间】:2019-10-09 03:51:20
【问题描述】:

我有一个在 Java 中运行的 ActiveMQ 使用者脚本,我在 while(true) 循环中调用 consumer.receive()

我需要为处理的每条消息实现超时(例如:如果一条消息处理超过 15 秒,我必须接收下一条)。

我已经为 ACK 提供了客户端确认模式。

请看我实现消费的consumeMessage方法。

期望的结果:

15 秒后需要丢弃第一条消息(即它不应调用acknowledge())。而是需要处理下一条消息。

//package consumer;
import org.apache.activemq.ActiveMQConnectionFactory;

import javax.jms.Connection;
import javax.jms.DeliveryMode;
import javax.jms.Destination;
import javax.jms.ExceptionListener;
import javax.jms.JMSException;
import javax.jms.Message;
import javax.jms.MessageConsumer;
import javax.jms.MessageProducer;
import javax.jms.Session;
import javax.jms.TextMessage;

public class ActivmqConsumer implements ExceptionListener {

    ActiveMQConnectionFactory connectionFactory = null;
    Connection connection = null;
    Session session = null;

    public ActivmqConsumer() throws Exception{
        String USERNAME = "admin";      
        String PASSWORD = "admin";
        this.connectionFactory = new ActiveMQConnectionFactory(USERNAME, PASSWORD, "tcp://192.168.56.101:61616?jms.prefetchPolicy.all=1");
        // Create a Connection
        this.connection = connectionFactory.createConnection();
        connection.start();
        connection.setExceptionListener(this);
        // Create a Session
        this.session = connection.createSession(false, Session.CLIENT_ACKNOWLEDGE);
    }

    public void consumeMessage(String destinationName, EventProcesser eventprocess){
        Destination destination = null;
        MessageConsumer consumer = null;
        try{
        // Create the destination (Topic or Queue)
        destination = session.createQueue(destinationName);

        // Create a MessageConsumer from the Session to the Topic or Queue
        consumer = session.createConsumer(destination);

        // Wait for a message
        while(true){
            Message message = consumer.receive(2);
            if(message==null){
                continue;
            }
            else if(message instanceof TextMessage) {
                TextMessage textMessage = (TextMessage) message;
                String text = textMessage.getText();
                System.out.println("Received: " + text);
                eventprocess.processEvent(text);
                message.acknowledge();
            } else{
                System.out.println("Received: " + message);
            }
         }
        }catch(Exception ex){
            ex.printStackTrace();
        }finally{
            try{
            consumer.close();
            }catch(Exception ex){
                ex.printStackTrace();
            }
        }
    }
}

【问题讨论】:

  • 我的回答是否解决了您的问题?如果是这样,请将其标记为正确,以帮助将来有相同问题的其他用户。如果不是,请详细说明原因。谢谢!

标签: java jms activemq amqp


【解决方案1】:

ActiveMQ 中没有“最大消息处理时间”或等效功能。您需要自己监控处理过程。也许看看this question/answer 了解如何做到这一点。另一种方法是使用 JTA 事务管理器并在超时为 15 秒的事务中使用消息。在 Java EE 容器中使用 MDB 是获得事务超时功能的一种简单方法。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2019-10-03
    • 1970-01-01
    • 2014-10-09
    • 2017-12-07
    • 2016-06-08
    • 2011-04-15
    • 2021-03-06
    • 2015-11-25
    相关资源
    最近更新 更多