【问题标题】:Is there any way to close MqttClient Threads spawned in the backend?有什么方法可以关闭后端产生的 MqttClient 线程?
【发布时间】:2019-01-09 10:12:27
【问题描述】:

我们有一个流式应用程序从MQTT 获取数据并加载到其他资源中。而且这个应用程序有多个线程来处理一些任务。

这里我们有两个任务(线程):

  • 第一个是阅读器
  • 第二个是作家

所以 READER 将从 MQTT 代理读取数据并写入 java 队列,WRITER 将从该队列中获取这些数据并将其写入一个数据库。此应用程序本身监视这些线程以查找任何故障。如果任何一个线程失败,那么我们将优雅地停止剩余的线程。在 paho MqttClient 类(READER 类)的情况下,即使它是一个线程类,也不会创建线程。但它会在后台创建多个线程。

因此,我们无法通过 java isAlive() 函数检查这些线程是否失败或正在运行。所以我们只是通过 MqttClient isConnected() 方法检查这个类是否有连接。一旦 isConnected 方法返回 false (5 次) ,那么我们将优雅地停止 Writer 线程。但是在后台产生的 Reader 类线程无法停止。我试过disconnect()close() 方法。但它并没有停止任何后台线程。它抛出错误断开的线程无法停止。

所以请大家帮忙。

【问题讨论】:

    标签: java mqtt paho


    【解决方案1】:

    您的建议听起来像是一个尴尬的设计。

    为什么不直接使用 Paho 回调,尤其是下面的 connectionLost

    private final MqttCallbackExtended mCallback = new MqttCallbackExtended() {
        @Override
        public void connectComplete(boolean reconnect, String brokerAddress) {
                mqttClient.subscribe("topic", 1, null, mSubscribeCallback);
        }
    
        @Override
        public void connectionLost(Throwable ex) {
        }
    
        @Override
        public void deliveryComplete(IMqttDeliveryToken deliveryToken) {
        }
    
        @Override
        public void messageArrived(String topic, MqttMessage mqttMessage) throws Exception {
        }
    };
    
    private final IMqttActionListener mConnectionCallback = new IMqttActionListener() {
        @Override
        public void onSuccess(IMqttToken asyncActionToken) {
            // do nothing, this case is handled in mCallback.connectComplete()
        }
    
        @Override
        public void onFailure(IMqttToken asyncActionToken, Throwable exception) {
        }
    };
    
    private final IMqttActionListener mSubscribeCallback = new IMqttActionListener() {
        @Override
        public void onSuccess(IMqttToken subscribeToken) {
        }
    
        @Override
        public void onFailure(IMqttToken subscribeToken, Throwable ex) {
        }
    
    };
    
    MqttConnectOptions connectOptions = new MqttConnectOptions();
    connectOptions.setCleanSession(true);
    connectOptions.setAutomaticReconnect(true);     
    connectOptions.setUserName("username");
    connectOptions.setPassword("password".toCharArray());
    
    MqttAsyncClient mqttClient = new MqttAsyncClient("tcp:// test.mosquitto.org");
    mqttClient.setCallback(mCallback);
    
    try {
        mqttClient.connect(connectOptions, null, mConnectionCallback);
    
    } catch (Exception ex) {
        System.err.println(ex.toString());
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2012-02-22
      • 1970-01-01
      • 2015-02-20
      • 2016-05-28
      • 2021-11-10
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多