From ba1e219d97e123b91657dc424f13a2fdeaa85d21 Mon Sep 17 00:00:00 2001
From: ‘liusuyi’ <1951119284@qq.com>
Date: Mon, 21 Aug 2023 13:37:01 +0800
Subject: [PATCH] 修改兴趣点列表井号模糊查询

---
 ard-work/src/main/java/com/ruoyi/utils/mqtt/MqttConsumerCallback.java |   22 ++++++++++++++--------
 1 files changed, 14 insertions(+), 8 deletions(-)

diff --git a/ard-work/src/main/java/com/ruoyi/utils/mqtt/MqttConsumerCallback.java b/ard-work/src/main/java/com/ruoyi/utils/mqtt/MqttConsumerCallback.java
index aa334a5..09fa6c9 100644
--- a/ard-work/src/main/java/com/ruoyi/utils/mqtt/MqttConsumerCallback.java
+++ b/ard-work/src/main/java/com/ruoyi/utils/mqtt/MqttConsumerCallback.java
@@ -1,7 +1,8 @@
 package com.ruoyi.utils.mqtt;
 
-import com.ruoyi.alarm.globalAlarm.service.impl.GlobalAlarmServiceImpl;
+import com.ruoyi.alarm.global.service.impl.GlobalAlarmServiceImpl;
 import com.ruoyi.common.utils.spring.SpringUtils;
+import com.ruoyi.storage.minio.service.IStorageMinioEventService;
 import lombok.extern.slf4j.Slf4j;
 import org.eclipse.paho.client.mqttv3.*;
 
@@ -36,9 +37,9 @@
     @Override
     public void connectionLost(Throwable cause) {
         log.info("MQTT连接断开,发起重连......");
-        try {
-            while (!client.isConnected()) {
-                Thread.sleep(5000);
+        while (!client.isConnected()) {
+            try {
+                Thread.sleep(10000);
                 if (null != client && !client.isConnected()) {
                     client.reconnect();
                     log.error("尝试重新连接");
@@ -46,9 +47,9 @@
                     client.connect(options);
                     log.error("尝试建立新连接");
                 }
+            } catch (Exception e) {
+                log.error("断开重连异常:" + e.getMessage());
             }
-        } catch (Exception e) {
-            e.printStackTrace();
         }
     }
 
@@ -68,12 +69,17 @@
     public void messageArrived(String topic, MqttMessage message) {
         try {
             // subscribe后得到的消息会执行到这里面
-            log.info("接收消息 【主题】:" + topic + " 【内容】:" + new String(message.getPayload()));
+            log.debug("接收消息 【主题】:" + topic + " 【内容】:" + new String(message.getPayload(), StandardCharsets.UTF_8));
             //进行业务处理(接收报警数据)
             GlobalAlarmServiceImpl globalAlarmService = SpringUtils.getBean(GlobalAlarmServiceImpl.class);
             globalAlarmService.receiveAlarm(topic, new String(message.getPayload(), StandardCharsets.UTF_8));
+            if (topic.equals("minioEvent"))
+            {
+                IStorageMinioEventService storageMinioEventService = SpringUtils.getBean(IStorageMinioEventService.class);
+                storageMinioEventService.parseStorageMinioEvent(new String(message.getPayload(), StandardCharsets.UTF_8));
+            }
         } catch (Exception e) {
-            log.info("处理mqtt消息异常:" + e);
+            log.debug("处理mqtt消息异常:" + e);
         }
     }
 

--
Gitblit v1.9.3