【问题标题】:IBM MQ How do I know I have read all topics?IBM MQ 我如何知道我已阅读所有主题?
【发布时间】:2017-06-02 05:26:37
【问题描述】:

我有一个动态主题树,每个主题的最后一条消息都“保留”在其中。当我订阅主题树(使用 JMS/Java/MQLibs)时,由于它是动态的并且我使用通配符“#”订阅,我无法知道我会提前收到哪些主题。所以,

  1. 我怎样才能知道我已经阅读了所有可用主题至少一次?

  2. 我如何知道阅读的主题是订阅前存在的原始“保留”消息还是订阅后的更新? (我是否需要根据主题保留 MessageID 的本地映射?我假设我无法将订阅时的时间戳与消息创建时间进行比较,因为服务器和客户端上的时钟可能不同)。

我添加了一些示例代码,以表明我访问的是主题而不是队列。

public class JMSServerSubscriber implements MessageListener { 

public JMSServerSubscriber() throws JMSException { 
    TopicConnection topicCon = JmsTopicConnectionFactory.createTopicConnection(); 
    TopicSession topicSes = topicCon.createTopicSession(false, Session.AUTO_ACKNOWLEDGE); 
    Topic topic = topicSes.createTopic("#"); 
    TopicSubscriber topicSub = topicSes.createSubscriber(topic); 
    topicSub.setMessageListener(this); 
    topicCon.start(); 
} 

    @Override 
    public void onMessage(Message arg0) { 
            BytesMessage bm = (BytesMessage) arg0; 
            try { 
                    String destination = bm.getJMSDestination().toString(); 
            } catch (JMSException e) { 
                    e.printStackTrace(); 
            }  
    } 
} 

【问题讨论】:

  • 我不明白这个问题。当您订阅一个主题时,您指定一个队列,MQ 会将发布到该主题的消息放入该队列。当该队列为空时,您已阅读发布到该主题的所有消息。我看不出订阅主题树的一部分有何不同。
  • 让我们说实际的主题树看起来像这样; /BusRoute/1/BusRoute/2,我订阅了/BusRoute/#。我会立即收到12 的回复,但我怎么知道这就是全部内容?

标签: java jms ibm-mq mq


【解决方案1】:

当没有更多发布时,consumer.receive() 方法将抛出带有 MQRC 2033 原因码的异常。您可以使用原因码确认您已阅读所有出版物。

在收到消息时调用 getIntProperty(JmsConstants.JMS_IBM_RETAIN);了解出版物是否保留。

【讨论】:

  • 您好,感谢您抽出宝贵时间回答。我虽然 2033 用于队列访问,但我访问主题(我知道主题实际上是内部队列)。我已经编辑了我的问题以包含示例代码。
  • 订阅将附加一个队列,MQ 将发布发布到该队列,并且您的应用程序从该队列接收消息。如果使用消息监听器,只有当消息到达订阅队列时才会调用 onMessage 方法。如果没有消息,则不会调用 onMessage 方法。
  • 感谢您的回复。我查看了 Message 的标题,并看到了描述底层队列的 XMSC_WMQ_BROKER 变量。但是,当我打开多个具有相同订阅的 JVM 时,每个 JVM 中的内部队列名称都是相同的。那么如何识别我的唯一队列,或共享队列上的个人订阅更新?
  • 您无需担心内部使用的是什么队列。 MQ JMS 客户端将选择用于您的客户端应用程序的消息。由于您使用的是“#”,因此您将获得所有出版物。
  • 是的,但我想知道我已经阅读了我的订阅允许的每个主题的一个示例。我应该求助于 PCF 并使用 DISPLAY TPSTATUS ('myPath/#') 进行查询吗?有没有 JMS 方法来实现这一点?
【解决方案2】:

这是我使用 PCF 的解决方案。请原谅任何语法错误。

public class TopicChecker implements MessageListener { 

    PCFMessageAgent           agent; 
    PCFMessage                pcfm; 
    Set<String>               allTopics; 
    TopicConnection           topicCon; 
    TopicSession              topicSes; 
    Topic                     topic; 
    TopicSubscriber           topicSub; 
    JmsTopicConnectionFactory jcf; 
    int                       total; 

    public TopicChecker() throws Exception { 
        allTopics = new HashSet<>(); 
        MQEnvironment.hostname = "myMqServer"; 
        MQEnvironment.channel = "MY.CHANNEL"; 
        MQEnvironment.port = 1515; 
        MQQueueManager m = new MQQueueManager("MY.QUEUE.MANAGER"); 
        agent = new PCFMessageAgent(m); 
        pcfm = new PCFMessage(MQConstants.MQCMD_INQUIRE_TOPIC_STATUS); 
        pcfm.addParameter(MQConstants.MQCA_TOPIC_STRING, "MyRoot/#"); 
        PCFMessage[] responses = agent.send(pcfm); 

        for (PCFMessage response: responses) { 
            /* We only publish to leaf nodes, so ignore branches. */
            if (response.getIntParameterValue(MQConstants.MQIACF_RETAINED_PUBLICATION) > 0) { 
            allTopics.add(response.getStringParameterValue(MQConstants.MQCA_TOPIC_STRING)); 
            } 
        } 

        total = allTopics.size(); 
        agent.disconnect(); 

        jcf = ...
        topicCon = jcf.createTopicConnection(); 
        topicSes = topicCon.createTopicSession(false, Session.AUTO_ACKNOWLEDGE); 
        topic = topicSes.createTopic("MyRoot/#"); 
        topicSub = topicSes.createSubscriber(topic); 
        topicSub.setMessageListener(this); 
        topicCon.start(); 
    } 

    @Override 
    public void onMessage(Message m) { 
     try { 
        String topicString = m.getJMSDestination().toString().replaceAll("topic://", ""); 
        allTopics.remove(topicString); 
        System.out.println("Read : " + topicString + " " + allTopics.size() + " of " + total + " remaining."); 
        if (allTopics.size() == 0) System.out.println("---------------------DONE----------------------"); 
        } catch (JMSException e) { 
               e.printStackTrace(); 
        } 
   }


   public static void main(String[] args) throws Exception { 
       TopicChecker tc = new TopicChecker(); 
       while (tc.allTopics.size() != 0); 
   } 
} 

【讨论】:

    【解决方案3】:

    if (response.getIntParameterValue(MQConstants.MQIACF_RETAINED_PUBLICATION) > 0) { allTopics.add(response.getStringParameterValue(MQConstants.MQCA_TOPIC_STRING)); }

    从上面的状态响应中,要读取主题名称,您实际上应该使用“MQCA_ADMIN_TOPIC_NAME”。

    【讨论】:

    • 欢迎来到 Stack Overflow!请正确格式化代码,缩进4个空格,并添加详细说明。
    猜你喜欢
    • 1970-01-01
    • 2015-04-05
    • 2018-07-03
    • 1970-01-01
    • 1970-01-01
    • 2011-10-07
    • 1970-01-01
    • 2017-01-11
    • 1970-01-01
    相关资源
    最近更新 更多