From 157c26f5188c7ed62a4547f7e3b5a5a3e3ed7729 Mon Sep 17 00:00:00 2001
From: ‘liusuyi’ <1951119284@qq.com>
Date: Sun, 08 Oct 2023 09:40:47 +0800
Subject: [PATCH] 优化mqtt生产者取消消费者订阅

---
 src/main/java/com/ard/utils/mqtt/MqttProducer.java |   65 ++++++--------------------------
 1 files changed, 13 insertions(+), 52 deletions(-)

diff --git a/src/main/java/com/ard/utils/mqtt/MqttConsumer.java b/src/main/java/com/ard/utils/mqtt/MqttProducer.java
similarity index 66%
rename from src/main/java/com/ard/utils/mqtt/MqttConsumer.java
rename to src/main/java/com/ard/utils/mqtt/MqttProducer.java
index fc03b39..aef7240 100644
--- a/src/main/java/com/ard/utils/mqtt/MqttConsumer.java
+++ b/src/main/java/com/ard/utils/mqtt/MqttProducer.java
@@ -1,21 +1,19 @@
 package com.ard.utils.mqtt;
 
 import lombok.extern.slf4j.Slf4j;
-import org.apache.commons.lang3.StringUtils;
 import org.eclipse.paho.client.mqttv3.*;
 import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;
 import org.springframework.beans.factory.annotation.Value;
 import org.springframework.boot.ApplicationArguments;
 import org.springframework.boot.ApplicationRunner;
 import org.springframework.core.annotation.Order;
-import org.springframework.expression.spel.ast.NullLiteral;
 import org.springframework.stereotype.Component;
 
 import java.io.UnsupportedEncodingException;
 
 /**
- * @Description: mqtt消费客户端
- * @ClassName: MqttConsumer
+ * @Description: mqtt生产客户端
+ * @ClassName: MqttProducer
  * @Author: 刘苏义
  * @Date: 2023年05月29日9:55
  * @Version: 1.0
@@ -23,11 +21,9 @@
 @Component
 @Slf4j(topic = "mqtt")
 @Order(1)
-public class MqttConsumer implements ApplicationRunner {
+public class MqttProducer implements ApplicationRunner {
     @Value("${spring.mqtt.enabled}")
     private Boolean MQTT_ENABLED;
-    @Value("${spring.mqtt.topic}")
-    private String MQTT_TOPIC;
     @Value("${spring.mqtt.host}")
     private String MQTT_HOST;
     @Value("${spring.mqtt.clientId}")
@@ -60,11 +56,8 @@
             getClient();
             // 2 设置配置
             MqttConnectOptions options = getOptions();
-            String[] topic = MQTT_TOPIC.split(",");
-            // 3 消息发布质量
-            int[] qos = getQos(topic.length);
-            // 4 最后设置
-            create(options, topic, qos);
+            // 3 最后设置
+            create(options);
         } catch (Exception e) {
             log.error("mqtt连接异常:" + e);
         }
@@ -97,69 +90,37 @@
         // 设置会话心跳时间
         options.setKeepAliveInterval(MQTT_KEEP_ALIVE);
         // 是否清除session
-        options.setCleanSession(true);
+        options.setCleanSession(false);
         log.debug("--生成mqtt配置对象");
         return options;
     }
 
-    /**
-     * qos   --- 3 ---
-     */
-    public int[] getQos(int length) {
-
-        int[] qos = new int[length];
-        for (int i = 0; i < length; i++) {
-            /**
-             *  MQTT协议中有三种消息发布服务质量:
-             *
-             * QOS0: “至多一次”,消息发布完全依赖底层 TCP/IP 网络。会发生消息丢失或重复。这一级别可用于如下情况,环境传感器数据,丢失一次读记录无所谓,因为不久后还会有第二次发送。
-             * QOS1: “至少一次”,确保消息到达,但消息重复可能会发生。
-             * QOS2: “只有一次”,确保消息到达一次。这一级别可用于如下情况,在计费系统中,消息重复或丢失会导致不正确的结果,资源开销大
-             */
-            qos[i] = 1;
-        }
-        log.debug("--设置消息发布质量");
-        return qos;
-    }
 
     /**
-     * 装载各种实例和订阅主题  --- 4 ---
+     * 连接并装载回调 --- 3 ---
      */
-    public void create(MqttConnectOptions options, String[] topic, int[] qos) {
+    public void create(MqttConnectOptions options) {
         try {
-            client.setCallback(new MqttConsumerCallback(client, options, topic, qos));
+            client.setCallback(new MqttProducerCallback(client, options));
             log.debug("--添加回调处理类");
             client.connect(options);
         } catch (Exception e) {
-            log.info("装载实例或订阅主题异常:" + e);
+            log.info("连接并装载回调异常:" + e);
         }
     }
 
-    /**
-     * 订阅某个主题
-     *
-     * @param topic
-     * @param qos
-     */
-    public void subscribe(String topic, int qos) {
-        try {
-            log.debug("topic:" + topic);
-            client.subscribe(topic, qos);
-        } catch (MqttException e) {
-            e.printStackTrace();
-        }
-    }
+
 
     /**
      * 发布,非持久化
      * <p>
-     * qos根据文档设置为1
+     * qos根据文档设置为2
      *
      * @param topic
      * @param msg
      */
     public static void publish(String topic, String msg) {
-        publish(1, false, topic, msg);
+        publish(2, false, topic, msg);
     }
 
     /**

--
Gitblit v1.9.3