【问题标题】:How can I tell an ACK corresponds to which publish message on MQTT?如何判断 ACK 对应于 MQTT 上的哪个发布消息?
【发布时间】:2016-12-21 16:52:51
【问题描述】:

我正在努力使用 Mqtt paho 驱动程序...

每当收到我的发布时,我都会使用 IMqttDeliveryToken 从服务器获取确认。

为了将它与实际的发布消息进行比较,我在 MqttMessage 上设置了一个 ID,以便从 IMqttDeliveryToken 中检索它...但它不起作用... IMqttDeliveryToken.getMessageId() 返回一个不正确的 ID当我尝试在 IMqttDeliveryToken.getMessage() 之后获取 ID 时,当 QoS 不为 0 时,它返回一个 NPE。

阅读 Javadoc 后,我了解到这是通常的行为:

在邮件送达之前,正在送达的邮件将被退回。消息传递后,将返回 null。

这让我想到另一个问题...在 Broker 发送确认后真的调用了 deliveryComplete() 方法吗?

这是我的代码:

client.setCallback(new MqttCallback() {
    @Override
    public void connectionLost(Throwable thrwbl) { }

    @Override
    public void messageArrived(String string, MqttMessage mm) throws Exception { }

    @Override
    public void deliveryComplete(IMqttDeliveryToken token) {
        try {
            System.out.println("Message ID from getMessageId() method : " + token.getMessageId());
            MqttMessage message = token.getMessage();
            System.out.println("Message ID from getMessage() method : " + message.getId());
        } catch (MqttException ex) {
            System.out.println(ex);
        } catch (Exception ex) {
            System.out.println(ex);
        }
    }
});

MqttMessage message = new MqttMessage();
message.setId(76);
message.setPayload("pouet".getBytes());
message.setQos(0);

client.publish("TEST", message);

将 QoS 设为 0:

Message ID from getMessageId() method : 1
Message ID from getMessage() method : 76

将 QoS 设为 1:

Message ID from getMessageId() method : 1
java.lang.NullPointerException

【问题讨论】:

    标签: java mqtt paho


    【解决方案1】:

    正如 git MqttMessage.java 中提到的那样

    /**
     * This is only to be used internally to provide the MQTT id of a message
     * received from the server.  Has no effect when publishing messages.
     * @param messageId
     */
    public void setId(int messageId) {
        this.messageId = messageId;
    }
    

    这是在发布消息时没有使用的地方。现在要了解为什么来自 getMessageId() 方法的消息 ID:1 发生,请查看以下内容。

    public IMqttDeliveryToken publish(String topic, MqttMessage message, Object userContext, IMqttActionListener callback) throws MqttException,
                MqttPersistenceException {
            final String methodName = "publish";
            //@TRACE 111=< topic={0} message={1}userContext={1} callback={2}
            log.fine(CLASS_NAME,methodName,"111", new Object[] {topic, userContext, callback});
    
            //Checks if a topic is valid when publishing a message.
            MqttTopic.validate(topic, false/*wildcards NOT allowed*/);
    
            MqttDeliveryToken token = new MqttDeliveryToken(getClientId());
            token.setActionCallback(callback);
            token.setUserContext(userContext);
            token.setMessage(message);
            token.internalTok.setTopics(new String[] {topic});
    
            MqttPublish pubMsg = new MqttPublish(topic, message);
            comms.sendNoWait(pubMsg, token);
    
            //@TRACE 112=<
            log.fine(CLASS_NAME,methodName,"112");
    
            return token;
        }
    

    MqttDeliveryToken 此处没有设置消息ID。发布时创建MqttPublish实例,内部扩展多级为MqttWireMessage.java,默认值为0。

    public MqttWireMessage(byte type) {
            this.type = type;
            // Use zero as the default message ID.  Can't use -1, as that is serialized
            // as 65535, which would be a valid ID.
            this.msgId = 0;
        }
    

    当在 ClientState.java 中为 mqtt 发布调用 Final Send 时,msgId = 0 的 MqttWireMessage 实例(内部 MqttPublish 是从上面的代码发送)被转发,因此 if 条件为真,调用 getNextMessageId()返回 1(因为它是第一条消息,否则它会根据最后一个 msg id 返回后续值)并设置为您在 deliveryComplete() 中跟踪的代码中的令牌。

    public void send(MqttWireMessage message, MqttToken token) throws MqttException {
            final String methodName = "send";
            if (message.isMessageIdRequired() && (message.getMessageId() == 0)) {
                message.setMessageId(getNextMessageId());
            }
            if (token != null ) {
                try {
                    token.internalTok.setMessageID(message.getMessageId());
                } catch (Exception e) {
                }
            }
    
            /////......
        }
    

    【讨论】:

    • 感谢您的启发:) 我使用了一个IMqttActionCallback 传入发布函数的参数而不是MqttCallback。
    【解决方案2】:

    回答您的下一个问题:deliveryComplete() 方法是否真的在经纪人发送确认后调用 --- 是的!!

    从这小段代码调用已完成的交付。

    private void handleActionComplete(MqttToken token)
                throws MqttException {
            final String methodName = "handleActionComplete";
            synchronized (token) {
                // @TRACE 705=callback and notify for key={0}
                log.fine(CLASS_NAME, methodName, "705", new Object[] { token.internalTok.getKey() });
                if (token.isComplete()) {
                    // Finish by doing any post processing such as delete 
                    // from persistent store but only do so if the action
                    // is complete
                    clientState.notifyComplete(token);
                }
    
                // Unblock any waiters and if pending complete now set completed
                token.internalTok.notifyComplete();
    
                if (!token.internalTok.isNotified()) {
                    // If a callback is registered and delivery has finished 
                    // call delivery complete callback. 
                    if ( mqttCallback != null 
                        && token instanceof MqttDeliveryToken 
                        && token.isComplete()) {
                            mqttCallback.deliveryComplete((MqttDeliveryToken) token);
                    }
                    // Now call async action completion callbacks
                    fireActionEvent(token);
                }
    
                // Set notified so we don't tell the user again about this action.
                if ( token.isComplete() ){
                   if ( token instanceof MqttDeliveryToken || token.getActionCallback() instanceof IMqttActionListener ) {
                        token.internalTok.setNotified(true);
                    }
                }
    
    
    
            }
        }
    

    即一旦收到确认,即完成通知完成,设置标志,调用deliveryComplete方法。

    【讨论】:

      猜你喜欢
      • 2022-09-22
      • 2022-06-14
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2011-12-31
      相关资源
      最近更新 更多