【问题标题】:Issues getting ActiveMQ Advisory messages for MessageConsumed获取 MessageConsumed 的 ActiveMQ 咨询消息的问题
【发布时间】:2015-07-14 02:53:33
【问题描述】:

我需要能够在 ActiveMQ 客户端使用 MQTT 消息时接收通知。

activemq.xml

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

在下面的代码中,我在 myTopic 上收到了 MQTT 消息。我没有在 processAdvisoryMessage / processAdvisoryBytesMessage 中收到咨询消息。

@Component
public class MqttMessageListener {
    @JmsListener(destination = "mytopic")
    public void processMessage(BytesMessage message) {
    }

    @JmsListener(destination = "ActiveMQ.Advisory.MessageConsumed.Topic.>")
    public void processAdvisoryMessage(Message message) {
        System.out.println("processAdvisoryMessage Got a message");
    }

    @JmsListener(destination = "ActiveMQ.Advisory.MessageConsumed.Topic.>")
    public void processAdvisoryBytesMessage(BytesMessage message) {
        System.out.println("processAdvisoryBytesMessageGot a message");
    }
}

我做错了什么?

我也尝试过使用 ActiveMQ BrokerFilter:

public class AMQMessageBrokerFilter extends GenericBrokerFilter {
    @Override
    public void acknowledge(ConsumerBrokerExchange consumerExchange, MessageAck ack) throws Exception {
        super.acknowledge(consumerExchange, ack);
    }

@Override
public void postProcessDispatch(MessageDispatch messageDispatch) {
    Message message = messageDispatch.getMessage();
}

@Override
public void messageDelivered(ConnectionContext context, MessageReference messageReference) {
    log.debug("messageDelivered called.");
    super.messageDelivered(context, messageReference);
}

@Override
public void messageConsumed(ConnectionContext context, MessageReference messageReference) {
    log.debug("messageConsumed called.");
    super.messageConsumed(context, messageReference);
}   

在第二种情况下,我无法同时拥有消息和用于发送消费通知的联系人。 confirm/messageDelivered/messageConsumed 都有一个连接上下文,但只有 postProcessDispatch 有我需要其中一部分的消息(有效负载是 JSON),以便发送我的传出消息。我可能会急切地使用 send 两者都有,但至少等到它被确认后更安全。

我试过了:

@Override
public void postProcessDispatch(MessageDispatch messageDispatch) {
    super.postProcessDispatch(messageDispatch);
    String topic =  messageDispatch.getDestination().getPhysicalName();
    if( topic == null || topic.equals("delivered") )
        return;

    try {
        ActiveMQTopic responseTopic = new ActiveMQTopic("delivered");
        ActiveMQTextMessage responseMsg = new ActiveMQTextMessage();
        responseMsg.setPersistent(false);
        responseMsg.setResponseRequired(false);
        responseMsg.setProducerId(new ProducerId());
        responseMsg.setText("Delivered msg: "+msg);
        responseMsg.setDestination(responseTopic);
        String messageKey = ":"+rand.nextLong();
        MessageId msgId = new MessageId(messageKey);
        responseMsg.setMessageId(msgId);

        ProducerBrokerExchange producerExchange=new ProducerBrokerExchange();
        ConnectionContext context = getAdminConnectionContext();
        producerExchange.setConnectionContext(context);
        producerExchange.setMutable(true);
        producerExchange.setProducerState(new ProducerState(new ProducerInfo()));
        next.send(producerExchange, responseMsg); 
    } 
    catch (Exception e) {
        log.debug("Exception: "+e);
    }

但是以上似乎导致服务器不稳定。我认为这与使用似乎错误的 getAdminConnectionContext 有关。

【问题讨论】:

    标签: jms activemq spring-jms


    【解决方案1】:

    我的工厂默认将 setPubSubDomain 设置为 false。这会禁用主题咨询消息的连接。我将其设置为 true 并且事情开始起作用。请注意,队列不适用于此集合。为了解决这个问题,我创建了两个工厂并命名了它们的 bean。

        @Bean(name="main")
        public DefaultJmsListenerContainerFactory jmsListenerContainerFactory(ConnectionFactory connectionFactory) {
            DefaultJmsListenerContainerFactory factory =  new DefaultJmsListenerContainerFactory();
            factory.setConnectionFactory(connectionFactory);
    //        factory.setDestinationResolver(destinationResolver);
    //        factory.setPubSubDomain(true);
            factory.setConcurrency("3-10");
            return factory;
        }  
    

    【讨论】:

      猜你喜欢
      • 2019-01-15
      • 1970-01-01
      • 1970-01-01
      • 2016-10-23
      • 1970-01-01
      • 2017-11-06
      • 1970-01-01
      • 2012-12-06
      • 2018-02-03
      相关资源
      最近更新 更多