首页
学习
活动
专区
工具
TVP
发布
社区首页 >问答首页 >如何停止java方法,直到它收到来自MQTT代理的消息为止。

如何停止java方法,直到它收到来自MQTT代理的消息为止。
EN

Stack Overflow用户
提问于 2018-08-12 00:41:41
回答 1查看 1.1K关注 0票数 2

我有一个包含2个方法的类,每个方法分别发布和订阅一个公共的MQTT代理。

public class operations{
    public void fetch(){
    // do something

     mqttConn.mqttSubscriberFetch(MQTTtopic);

    // do something
   }

    public void store(){
   // do something

    mqttConn.mqttSubscriberStore(MQTTtopic);

   // do something
   }

}

方法fetch的MQTT方法如下:

public class mqttConnectionFetch implements  MqttCallback{

    MqttClient client;

    String topic;

    public mqttConnectionFetch() {
    }


    public void mqttSubscriberFetch(String topic) {

        final String MQTT_BROKER = "tcp://localhost:1883" ;
        final String CLIENT_ID = "fetch";
        try {
            client = new MqttClient(MQTT_BROKER, CLIENT_ID);
            client.connect();
            client.setCallback(this);
            client.subscribe(topic);
            MqttMessage message = new MqttMessage();

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



    @Override
    public void connectionLost(Throwable cause) {
        // TODO Auto-generated method stub

    }

    @Override
    public void messageArrived(String topic, MqttMessage message)
            throws Exception {

        System.out.println("the message received is "+message);

        if(message.toString().equals("Processed")) {
            MqttPublisherFetch("Processed", "operations/fetch");

        }

    }

    @Override
    public void deliveryComplete(IMqttDeliveryToken token) {
        // TODO Auto-generated method stub

    }


    public void MqttPublisherFetch(String message, String topic) {
        final String MQTT_BROKER = "tcp://localhost:1883" ;
        final String CLIENT_ID = "store";
        try {
            client = new MqttClient(MQTT_BROKER, CLIENT_ID);
            client.connect();
            createMqttMessage(message,topic);
            client.disconnect();
        } catch (MqttException e) {
            e.printStackTrace();
        }
    }

    private void createMqttMessage(String message, String topic) throws MqttException {
        MqttMessage publishMessage = new MqttMessage();
        publishMessage.setPayload(message.getBytes());
        client.publish(topic, publishMessage);

    }


}

现在我正在尝试这样一种功能,每当我的fetch方法订阅一个主题时,如果来自代理的消息是“正在处理”,它就应该进入等待状态。而store方法应该正常工作。当接收到的消息被“处理”时,抓取应该再次开始工作。

我用普通的java wait()和start()尝试了一下,但是我没有得到想要的输出。有人能帮我揭开这个神秘面纱吗?

EN

回答 1

Stack Overflow用户

发布于 2018-08-12 07:50:34

据我所知,store方法由几个步骤组成。其中一个步骤是通过MQTT发送消息并等待响应。从store方法的客户端的角度来看,它是同步的,即客户端只有在整个处理(包括通过MQTT进行异步发送/接收)完成之后才会收到响应。

这里你想要解决的是一个典型的问题,当一个线程需要等待,直到另一个线程发生了某些情况。有许多方法可以实现这一点,只需检查How to make a Java thread wait for another thread's output?中建议的选项数量。

最简单的方法是使用CountDownLatch。我将使用fetch方法,因为这就是您提供的代码。

首先,您将修改mqttSubscriberFetch以创建内部CountDownLatch对象:

private CountDownLatch processingFinishedLatch;

public void mqttSubscriberFetch(String topic) {

    final String MQTT_BROKER = "tcp://localhost:1883" ;
    final String CLIENT_ID = "fetch";
    try {
        client = new MqttClient(MQTT_BROKER, CLIENT_ID);
        client.connect();
        client.setCallback(this);
        client.subscribe(topic);
        MqttMessage message = new MqttMessage();
        processingFinishedLatch = new CountDownLatch(1);

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

则需要在消息接收回调中触发该信号:

@Override
public void messageArrived(String topic, MqttMessage message)
        throws Exception {

    System.out.println("the message received is "+message);

    if(message.toString().equals("Processed")) {
        MqttPublisherFetch("Processed", "operations/fetch");
        this.processingFinishedLatch.countDown();
    }

}

此外,您还需要为mqttConnectionFetch的客户端提供一种等待锁存器变为零的方法:

void waitForProcessingToFinish() throws InterruptedException {
     this.processingFinishedLatch.await();
}

fetch方法应该像这样使用它:

public void fetch(){
// do something

mqttConn.mqttSubscriberFetch(MQTTtopic);

// do everything that is needed to initiate processing

mqttConn.waitForProcessingToFinish()

// do something

}

执行fetch的线程将在waitForProcessingToFinish中等待,直到latch达到零,这将在适当的消息到来时发生。

当消息从未到达时,可以很容易地修改此方法来解决此问题。只需使用await with some timeout即可

boolean waitForProcessingToFinish(long timeout,
        TimeUnit unit) throws InterruptedException {
     return this.processingFinishedLatch.await();
}

fetch应该检查返回值,如果发生超时,可能会向调用者返回错误。

请注意,此解决方案有一个缺点,即正在处理fetch的线程将一直处于忙碌状态。如果这是用于处理传入HTTP请求的线程池中的线程,这可能会限制服务器可以并发处理的请求数量。有一些方法可以缓解这种情况,比如async servletsRxJavaSpring WebFlux。这个问题的实际解决方案在很大程度上取决于您用于REST API实现的技术。

票数 1
EN
页面原文内容由Stack Overflow提供。腾讯云小微IT领域专用引擎提供翻译支持
原文链接:

https://stackoverflow.com/questions/51801756

复制
相关文章

相似问题

领券
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档