【问题标题】:how do i subscribe to activemq to get notified when message has been processed by consumer我如何订阅activemq以在消费者处理消息时得到通知
【发布时间】:2015-10-05 11:52:39
【问题描述】:

我实际上正在寻找来自 ActiveMQ 的咨询或任何其他替代支持,以便在与 Consumer 关联的 MessageListener 完成处理消息时通知我。

MessageDelivered 建议似乎会在代理收到消息后立即通知。此外,MessageConsumed 建议声称在消费者收到消息时进行通知。

------------更新---------- --------

请在下面找到代码sn-p:

public class SampleListener implements MessageListener {

    private Session session;

    public SampleListener(Session session) {
        this.session = session; 
    }

    public void onMessage(Message message) {
        try {
             // do something
             session.commit();
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

public class SampleConsumer {

    private boolean stopConnection = false;

    public static void main(String[] args) {
        new SampleConsumer().start();
    }

    public void start() {
            ActiveMQConnectionFactory connectionFactory = new ActiveMQConnectionFactory("tcp://localhost:61616");
            Connection connection;
            try {
                connection = connectionFactory.createConnection();
                connection.start();
                Session session = connection.createSession(true, Session.SESSION_TRANSACTED);
                Destination destination = session.createTopic("test");
                MessageConsumer messageConsumer = session.createConsumer(destination);
                messageConsumer.setMessageListener(new SampleListener(session));

                try {
                    synchronized (this) {
                        while (!stopConnection) {
                            wait();
                        }
                    }
                } catch (InterruptedException e) {
                    e.printStackTrace();
                } finally {
                    session.close();
                    connection.close();
                }

            } catch (JMSException e) {
                e.printStackTrace();
            }
        }
    }

    public void stop() {
        synchronized (this) {
            stopConnection = true;
            notify();
        }
    }
}


public class SampleProducer implements MessageListener {

    private boolean messageDelivered;

    @Test
    public void shouldTestSomething() throws JMSException, InterruptedException {
        producerConnection = new ActiveMQConnectionFactory("tcp://localhost:61616").createConnection();
        producerConnection.start();
        Session session = producerConnection.createSession(true, SESSION_TRANSACTED);

        Destination destination = session.createTopic("test");
        MessageConsumer advisoryConsumer = session.createConsumer(AdvisorySupport.getMessageConsumedAdvisoryTopic(destination));
        advisoryConsumer.setMessageListener(this);

        Message message = session.createTextMessage("Hi");
        Destination destination = session.createTopic("test");
        MessageProducer producer = session.createProducer(destination);
        producer.setDeliveryMode(DeliveryMode.PERSISTENT);
        producer.send(message);
        session.commit();

        synchronized (this) {
            while (!messageDelivered) {
                wait();
            }
        }

        session.close();

        // some assertions
    }

    public void onMessage(Message message) {
        // do something

        synchronized (this) {
            messageDelivered = true;
            notify();
        }
    }
}

【问题讨论】:

  • 如果您将 ack-mode 设置为事务处理,MessageConsumed 咨询仍会提前发生,会发生什么情况?根据issues.apache.org/jira/browse/AMQ-3361的关闭原因,它不应该。
  • Aksel,ack-mode 已经被交易。虽然 MessageDelivered 咨询发生得很早,但我无法通过 MessageConsumed 咨询获得任何通知。我有 PolicyEntry 参数来打开 MessageConsumed,但仍然没有任何消息排队等待 MessageConsumed 咨询。
  • 我没有关注。问题是您认为消息消费会提前发生还是您无法为其启用建议?
  • 后者。我确实在管理控制台的主题部分中看到了 MessageConsumed 咨询,但此主题没有任何反应,即配置为收听此咨询的侦听器永远不会被调用。然而,在收听 MessageDelivered 咨询时,通知会在消费者完成消费消息之前发生。
  • 您在 onMessage() 中调用 session.commit() ?否则我无法猜测,并且需要查看一些代码才能更好地回答

标签: jms monitoring activemq


【解决方案1】:

默认情况下不启用某些建议。链接:

http://activemq.apache.org/advisory-message.html

可以通过将 policyEntry 添加到 activemq.xml 来启用已禁用的建议 http://activemq.apache.org/xml-configuration.html

将以下内容添加到activemq.xml:

     <destinationPolicy>
        <policyMap>
          <policyEntries>
            <policyEntry topic=">" advisoryForConsumed="true" />
            <policyEntry  queue=">" advisoryForConsumed="true" />
            ..            
          </policyEntries>
        </policyMap>
    </destinationPolicy>

在消费者中启用咨询并调用 session.commit() 后,咨询将被传递。

如果您使用嵌入式代理,您只需将 activemq.xml 放在类路径上并使用以下命令启动代理:

BrokerService broker = BrokerFactory.createBroker("xbean:activemq.xml",true);

(我没有找到任何方法在不使用 activemq.xml 的情况下启用禁用的建议)。

【讨论】:

  • 谢谢阿克塞尔。我一直在 policyEntry 中有advisoryForConsumed="true" (我使用外部代理,而不是嵌入式代理)。管理控制台确实显示了此建议的条目,并且确实反映了建议消息侦听器。但是没有建议消息在此建议中排队,因此没有消息到达侦听器。我的直觉是,这与生产者和消费者使用不同的会话以及处于不同的线程有关,但我不明白为什么这会是个问题。
  • 我让您的代码可以进行微小的更改,然后我粘贴它
  • 我猜我不会让你开心的,其中一个变化是使用队列而不是主题:消费者:pastebin.com/sJUqrVTi生产者:pastebin.com/XK0SGXnz配置:pastebin.com/zpN850st
  • 首先非常感谢您指出它适用于队列,我试过了,它也有效!幸运的是,我现在意识到我正在构建的应用程序宁愿在队列语义上工作,因为在订阅者注册之前,主题不会缓冲任何消息。所以这应该可以解决我的问题。但我仍然很好奇为什么 MessageConsumed 咨询不适用于主题。它是故意不工作的吗?还是我错过了什么?
  • 是否可以使用 PAHO 或其他 JS 客户端在 javascript 中订阅咨询主题
猜你喜欢
  • 1970-01-01
  • 2019-10-03
  • 2017-02-20
  • 1970-01-01
  • 1970-01-01
  • 2016-06-08
  • 2012-06-28
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多