最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

SpringBoot整合MQTT協(xié)議實現(xiàn)消息訂閱與發(fā)布功能

 更新時間:2025年09月10日 14:48:38   作者:燦燦不熬夜  
文章介紹了基于MQTT的項目實現(xiàn),包含依賴配置、啟動時固定主題訂閱、連接配置類、消息發(fā)布與訂閱功能、回調(diào)處理連接狀態(tài)及消息,并通過接口測試消息發(fā)送與接收,驗證EMQX連接狀態(tài)與主題交互,感興趣的朋友跟隨小編一起看看吧

1、相關(guān)依賴 pom.xml文件

<dependency>
    <groupId>org.springframework.integration</groupId>
    <artifactId>spring-integration-mqtt</artifactId>
</dependency>

2、配置文件 application.yml

這里的訂閱主題可不要,我這里用于啟動的時候就訂閱固定主題。適用主題固定的場景。

# MQTT服務(wù)地址,端口默認(rèn)1883
mqtt-broker-url: tcp://127.0.0.1:1883
# 用戶名
mqtt-username: admin
# 密碼
mqtt-password: public
# 訂閱主題(可以多個)
mqtt-default-topic: mqtt/topic_test
# 客戶端Id
mqtt-clientId: can

3、MQTT配置類

用于配置項目啟動時就連接MQTT。

@Component
public class MqttConfig {
    @Resource
    private MqttPushClient mqttPushClient;
    @Resource
    private MqttSubClient mqttSubClient;
    /**
     * 用戶名
     */
    @Value("${mqtt-username}")
    private String username;
    /**
     * 密碼
     */
    @Value("${mqtt-password}")
    private String password;
    /**
     * 連接地址
     */
    @Value("${mqtt-broker-url}")
    private String hostUrl;
    /**
     * 客戶Id
     */
    @Value("${mqtt-clientId}")
    private String clientId;
    /**
     * 默認(rèn)連接話題,多個的話用逗號隔開
     */
    @Value("${mqtt-default-topic}")
    private String defaultTopic;
    /**
     * 超時時間
     */
    private int timeout = 100;
    /**
     * 保持連接數(shù)
     */
    private int keepalive = 60;
    /**
     * 連接至mqtt服務(wù)器,獲取mqtt連接
     *
     * @return MqttPushClient
     */
    @Bean
    public MqttPushClient getMqttPushClient() {
        // 連接至mqtt服務(wù)器,獲取mqtt連接
        mqttPushClient.connect(hostUrl, clientId, username, password, timeout, keepalive);
        // 訂閱默認(rèn)主題
        mqttSubClient.subScribeDataPublishTopic(defaultTopic);
        return mqttPushClient;
    }
}

4、發(fā)布連接類

連接MQTT的方法、發(fā)布消息的方法。

@Slf4j
@Component
public class MqttPushClient {
    private static final Logger logger = LoggerFactory.getLogger(MqttPushClient.class);
    @Autowired
    private PushCallback pushCallback;
    @Getter
    private static MqttClient client;
    public static void setClient(MqttClient client) {
        MqttPushClient.client = client;
    }
    /**
     * 連接
     * @param host  mqtt://127.0.0.1:1883
     * @param clientId  can
     * @param username admin
     * @param password  password
     * @param timeout   100
     * @param keepalive 60
     */
    public void connect(String host, String clientId, String username, String password, int timeout, int keepalive) {
        MqttClient client;
        try {
            client = new MqttClient(host, clientId, new MemoryPersistence());
            MqttConnectOptions options = new MqttConnectOptions();
            options.setCleanSession(true);
            options.setUserName(username);
            options.setPassword(password.toCharArray());
            options.setConnectionTimeout(timeout);
            options.setKeepAliveInterval(keepalive);
            // automaticReconnect 為 true 表示斷線自動重連,但僅僅只是重新連接,并不訂閱主題;在 connectComplete 回調(diào)函數(shù)重新訂閱
            options.setAutomaticReconnect(true);
            MqttPushClient.setClient(client);
            try {
                //設(shè)置回調(diào)類
                client.setCallback(pushCallback);
                IMqttToken iMqttToken = client.connectWithResult(options);
                boolean complete = iMqttToken.isComplete();
                log.error("MQTT連接{}", complete ? "成功" : "失敗");
            } catch (Exception e) {
                logger.error(e.getMessage());
                e.printStackTrace();
            }
        } catch (Exception e) {
            logger.error(e.getMessage());
            e.printStackTrace();
        }
    }
    /**
     * 關(guān)閉MQTT連接
     */
    public void close() throws MqttException {
        client.disconnect();
        client.close();
    }
    /**
     * 發(fā)布,默認(rèn)qos為0,非持久化
     *
     * @param topic       主題名
     * @param pushMessage 消息
     */
    public void publish(String topic, String pushMessage) {
        publish(0, false, topic, pushMessage);
    }
    /**
     * 發(fā)布
     * QoS 0:消息最多傳送一次。如果當(dāng)前客戶端不可用,它將丟失這條消息。
     * QoS 1:消息至少傳送一次。
     * QoS 2:消息只傳送一次。
     * @param qos
     * @param retained
     * @param topic
     * @param pushMessage
     */
    public void publish(int qos, boolean retained, String topic, String pushMessage) {
        MqttMessage message = new MqttMessage();
        message.setQos(qos);
        message.setRetained(retained);
        message.setPayload(pushMessage.getBytes());
        MqttTopic mTopic = MqttPushClient.getClient().getTopic(topic);
        // MQTT主題不存在
        if (null == mTopic) return;
        try {
            mTopic.publish(message);
        } catch (Exception e) {
            log.error("MQTT發(fā)送消息異常:", e);
            e.printStackTrace();
        }
    }
}

5、訂閱類

用于訂閱某個或多個主題、取消訂閱某個或者多個主題。

@Slf4j
@Component
public class MqttSubClient {
    private static final Logger logger = LoggerFactory.getLogger(MqttSubClient.class);
    // 訂閱多個主題以逗號分開
    public void subScribeDataPublishTopic(String defaultTopic) {
        //訂閱test_queue主題
        String[] mqttTopic = defaultTopic.split(",");
        for (String s : mqttTopic) {
            //訂閱主題
            subscribe(s, 0);
        }
    }
    /**
     * 訂閱某個主題,qos默認(rèn)為0
     *
     * @param topic 主題
     */
    public void subscribe(String topic) {
        subscribe(topic, 0);
    }
    /**
     * 訂閱某個主題
     *
     * @param topic 主題名
     * @param qos   qos
     */
    public void subscribe(String topic, int qos) {
        try {
            MqttClient client = MqttPushClient.getClient();
            if (client == null) {
                return;
            }
            client.subscribe(topic, qos);
            log.error("MQTT訂閱主題:{}", topic);
        } catch (MqttException e) {
            logger.error(e.getMessage());
            e.printStackTrace();
        }
    }
    /**
     * 取消訂閱某個主題
     * @param topic 要取消訂閱的主題名
     */
    public void unsubscribe(String topic) {
        try {
            MqttClient client = MqttPushClient.getClient();
            if (client == null || !client.isConnected()) {
                return;
            }
            client.unsubscribe(topic); // 取消訂閱
            log.error("MQTT取消訂閱主題: {}", topic);
        } catch (MqttException e) {
            log.error("取消訂閱失敗: {}", e.getMessage());
            e.printStackTrace();
        }
    }
    /**
     * 批量取消訂閱多個主題
     * @param topics 主題數(shù)組
     */
    public void unsubscribe(String[] topics) {
        try {
            MqttClient client = MqttPushClient.getClient();
            if (client == null || !client.isConnected()) {
                return;
            }
            client.unsubscribe(topics); // 取消訂閱多個主題
            log.error("MQTT取消訂閱主題: {}", Arrays.toString(topics));
        } catch (MqttException e) {
            log.error("取消訂閱失敗: {}", e.getMessage());
            e.printStackTrace();
        }
    }
}
 

6、回調(diào)類

處理MQTT連接斷開重連、訂閱主題接收的消息處理。

@Slf4j
@Component
public class PushCallback implements MqttCallback {
    @Resource
    @Lazy
    private MqttPushClient mqttPushClient;
    /**
     * 連接丟失后,一般在這里面進行重連(重連的邏輯需要自己處理)
     * @param cause .
     */
    @Override
    public void connectionLost(Throwable cause) {
        log.error("MQTT連接斷開,正在重連:" + cause);
    }
    /**
     * 發(fā)送消息,消息到達后處理方法
     * @param token .
     */
    @Override
    public void deliveryComplete(IMqttDeliveryToken token) {
        log.error("deliveryComplete---------{}", token.isComplete());
    }
    /**
     * 訂閱主題接收到消息處理方法
     * @param topic 主題
     * @param message   消息
     */
    @Override
    public void messageArrived(String topic, MqttMessage message) {
        // 訂閱主題后得到的消息會執(zhí)行到這里面,這里在控制臺有輸出
        log.error("MQTT接收消息主題 : {}", topic);
        log.error("MQTT接收消息Qos : {}", message.getQos());
        log.error("MQTT接收消息內(nèi)容 : {}", message);
    }
}

7、啟動后,進入EMQX管理頁面

程序允許打印連接成功,去EMQX管理頁面查看。

EMQX管理頁面這里有所有的主題列表。

包括客戶端訂閱的主題。

8、通過接口給主題發(fā)送消息

@RestController
@Slf4j
@RequestMapping("/api")
public class ApiController {
    @Resource
    private MqttPushClient mqttPushClient;
    @GetMapping("/test")
    public String getVersions(@RequestParam String topic, @RequestParam String message) {
        mqttPushClient.publish(topic, message);
        return "ok";
    }
}

瀏覽器直接調(diào)用,topic:配置文件里面訂閱的主題。message:你想發(fā)送給主題的消息。

控制臺日志打?。喊l(fā)送成功,并且接收到了主題發(fā)送的消息。

到此這篇關(guān)于SpringBoot整合MQTT協(xié)議實現(xiàn)消息訂閱與發(fā)布功能的文章就介紹到這了,更多相關(guān)SpringBoot整合MQTT訂閱與發(fā)布內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Mybatis下劃線駝峰處理的幾種方法

    Mybatis下劃線駝峰處理的幾種方法

    這篇文章主要講述Mybatis下劃線駝峰處理的幾種方法,對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2023-12-12
  • 詳解如何在spring boot中使用spring security防止CSRF攻擊

    詳解如何在spring boot中使用spring security防止CSRF攻擊

    這篇文章主要介紹了詳解如何在spring boot中使用spring security防止CSRF攻擊,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2018-05-05
  • Java內(nèi)部類_動力節(jié)點Java學(xué)院整理

    Java內(nèi)部類_動力節(jié)點Java學(xué)院整理

    內(nèi)部類是指在一個外部類的內(nèi)部再定義一個類。下面通過本文給大家java內(nèi)部類的使用小結(jié),需要的朋友參考下吧
    2017-04-04
  • Java使用反射和動態(tài)代理實現(xiàn)一個View注解綁定庫

    Java使用反射和動態(tài)代理實現(xiàn)一個View注解綁定庫

    這篇文章主要介紹了Java使用反射和動態(tài)代理實現(xiàn)一個View注解綁定庫,代碼簡潔,使用簡單,擴展性強,結(jié)合實例代碼給大家介紹的非常詳細,需要的朋友可以參考下
    2022-05-05
  • java獲取兩個數(shù)組中不同數(shù)據(jù)的方法

    java獲取兩個數(shù)組中不同數(shù)據(jù)的方法

    這篇文章主要介紹了java獲取兩個數(shù)組中不同數(shù)據(jù)的方法,實例分析了java操作數(shù)組的技巧,非常具有實用價值,需要的朋友可以參考下
    2015-03-03
  • Java接口回調(diào)和方法回調(diào)的簡單實現(xiàn)步驟

    Java接口回調(diào)和方法回調(diào)的簡單實現(xiàn)步驟

    這篇文章主要介紹了Java接口回調(diào)和方法回調(diào)的相關(guān)資料,接口回調(diào)是一種設(shè)計模式,實現(xiàn)三方解耦,調(diào)用者提供接口實現(xiàn),文中通過代碼介紹的非常詳細,需要的朋友可以參考下
    2025-03-03
  • Spring?Boot?JWT登錄授權(quán)使用(無感刷新)

    Spring?Boot?JWT登錄授權(quán)使用(無感刷新)

    JWT作為一種輕量級的身份認(rèn)證與授權(quán)方案,憑借其無狀態(tài)、可跨域、易于擴展的特性,成為?Spring?Boot?項目中實現(xiàn)認(rèn)證授權(quán)的主流選擇,本文詳解Spring?Boot整合?JWT?實現(xiàn)登錄認(rèn)證與授權(quán)的全流程,助力開發(fā)者快速搭建安全可靠的認(rèn)證授權(quán)體系,感興趣的可以了解一下
    2026-04-04
  • SpringMVC之RequestContextHolder詳細解析

    SpringMVC之RequestContextHolder詳細解析

    這篇文章主要介紹了SpringMVC之RequestContextHolder詳細解析,正常來說在service層是沒有request的,然而直接從controlller傳過來的話解決方法太粗暴,后來發(fā)現(xiàn)了SpringMVC提供的RequestContextHolder,需要的朋友可以參考下
    2023-11-11
  • JAVA基礎(chǔ)之繼承(inheritance)詳解

    JAVA基礎(chǔ)之繼承(inheritance)詳解

    繼承(inheritance)是Java OOP中一個非常重要的概念。這篇文章主要介紹了JAVA基礎(chǔ)之繼承(inheritance),需要的朋友可以參考下
    2017-03-03
  • Java 創(chuàng)建兩個線程模擬對話并交替輸出實現(xiàn)解析

    Java 創(chuàng)建兩個線程模擬對話并交替輸出實現(xiàn)解析

    這篇文章主要介紹了Java 創(chuàng)建兩個線程模擬對話并交替輸出實現(xiàn)解析,文中通過示例代碼介紹的非常詳細,對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友可以參考下
    2019-10-10

最新評論

从江县| 新龙县| 沅陵县| 台州市| 吉隆县| 虹口区| 镇沅| 乌审旗| 阳东县| 青岛市| 德昌县| 牡丹江市| 文成县| 琼结县| 碌曲县| 兰州市| 会东县| 拉萨市| 金阳县| 沙河市| 宁陕县| 安乡县| 永和县| 盐源县| 五寨县| 洛宁县| 如东县| 荔波县| 湘西| 绥德县| 庆云县| 尼勒克县| 鹿邑县| 赣榆县| 宁海县| 靖安县| 井陉县| 库车县| 九江县| 榕江县| 江北区|