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

Springboot整合mqtt實(shí)現(xiàn)軟硬件通信的示例代碼

 更新時(shí)間:2025年12月26日 10:37:35   作者:程序陸  
本文介紹了如何使用Spring Boot整合MQTT協(xié)議,實(shí)現(xiàn)物聯(lián)網(wǎng)軟硬件通信, MQTT是一種輕量級(jí)的消息隊(duì)列遙測(cè)傳輸協(xié)議,適用于物聯(lián)網(wǎng)設(shè)備通信和遠(yuǎn)程監(jiān)控,本文給大家介紹的非常詳細(xì),感興趣的朋友跟隨小編一起看看吧

前言

本文實(shí)現(xiàn)Springboot整合mqtt,更好地實(shí)現(xiàn)物聯(lián)網(wǎng)軟硬件通信。

一、mqtt是什么?

MQTT 是消息隊(duì)列遙測(cè)傳輸的縮寫,是一種輕量級(jí)、基于發(fā)布 / 訂閱模式的物聯(lián)網(wǎng)通信協(xié)議。

應(yīng)用場(chǎng)景:

  • 物聯(lián)網(wǎng)設(shè)備通信,比如智能家居、工業(yè)傳感器、穿戴設(shè)備。
  • 遠(yuǎn)程監(jiān)控與數(shù)據(jù)采集,例如環(huán)境監(jiān)測(cè)、設(shè)備狀態(tài)上報(bào)。

注:這里實(shí)現(xiàn)的是軟件后端層面的mqtt,是對(duì)硬件消息的接收和發(fā)送,如果要做一個(gè)完整的物聯(lián)網(wǎng)項(xiàng)目,也需要硬件層面連接mqtt服務(wù)器,并且訂閱、發(fā)布相關(guān)主題的信息。

二、物聯(lián)網(wǎng)項(xiàng)目結(jié)構(gòu)圖

該圖為完整的mqtt項(xiàng)目結(jié)構(gòu)圖。其中mqtt服務(wù)器部分可以用官方提供的公共的服務(wù)器,也可以在自己的服務(wù)器上搭建,只需要在后端配置文件添加mqtt服務(wù)器ip地址等相關(guān)信息。

注:Springboot項(xiàng)目不是自己內(nèi)部開一個(gè)mqtt服務(wù),所以我們只需要連接并使用相應(yīng)的mqtt服務(wù)器即可。

三、實(shí)現(xiàn)代碼

在pom.xml中引入Maven坐標(biāo):

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

在Springboot配置文件中配置mqtt服務(wù)器相關(guān)信息

spring:
  mqtt:
    url: tcp://broker.emqx.io:1883
    #用戶名
    username: admin
    #密碼
    password: 123456
    #客戶端id(不能重復(fù))
    client:
      id: consumer-id
    #MQTT默認(rèn)的消息推送主題,實(shí)際可在調(diào)用接口時(shí)指定
    default:
      topic: topic

這里的mqtt服務(wù)器地址 tcp://broker.emqx.io:1883 是EMQX官方提供的一個(gè)公共的服務(wù)器,如果是想看一下初步效果,可以使用該服務(wù)器。

接下來我們可以放兩部分代碼,一部分就是接收消息,也就是訂閱消息的部分,另一部分就是用來發(fā)布相關(guān)消息。

訂閱消息部分:

@Configuration
public class MqttConsumerConfig {
    @Value("${spring.mqtt.username}")
    private String username;
    @Value("${spring.mqtt.password}")
    private String password;
    @Value("${spring.mqtt.url}")
    private String hostUrl;
    @Value("${spring.mqtt.client.id}")
    private String clientId;
    @Value("${spring.mqtt.default.topic}")
    private String defaultTopic;
    /**
     * 客戶端對(duì)象
     */
    private MqttClient client;
    /**
     * 在bean初始化后連接到服務(wù)器
     */
    @PostConstruct
    public void init(){
        connect();
    }
    /**
     * 客戶端連接服務(wù)端
     */
    public void connect(){
        try {
            client = new MqttClient(hostUrl,clientId,new MemoryPersistence());
            MqttConnectOptions options = new MqttConnectOptions();
            options.setCleanSession(true);
//            //設(shè)置連接用戶名
//            options.setUserName(username);
//            //設(shè)置連接密碼
//            options.setPassword(password.toCharArray());
            //設(shè)置超時(shí)時(shí)間,單位為秒
            options.setConnectionTimeout(100);
            //設(shè)置心跳時(shí)間 單位為秒,表示服務(wù)器每隔1.5*20秒的時(shí)間向客戶端發(fā)送心跳判斷客戶端是否在線
            options.setKeepAliveInterval(20);
            //設(shè)置遺囑消息的話題,若客戶端和服務(wù)器之間的連接意外斷開,服務(wù)器將發(fā)布客戶端的遺囑信息
            options.setWill("willTopic",(clientId + "與服務(wù)器斷開連接").getBytes(),0,false);
            //設(shè)置回調(diào)
            client.setCallback(new MqttConsumerCallBack());
            client.connect(options);
            //訂閱主題
            //消息等級(jí),和主題數(shù)組一一對(duì)應(yīng),服務(wù)端將按照指定等級(jí)給訂閱了主題的客戶端推送消息
            int[] qos = {1};
            System.out.println("連接");
            //主題
            //String[] topics = {"data","status"};
            client.subscribe(topics,qos);
        } catch (MqttException e) {
            e.printStackTrace();
        }
    }
    /**
     * 斷開連接
     */
    public void disConnect(){
        try {
            client.disconnect();
        } catch (MqttException e) {
            e.printStackTrace();
        }
    }
    /**
     * 訂閱主題
     */
    public void subscribe(String topic,int qos){
        try {
            client.subscribe(topic,qos);
        } catch (MqttException e) {
            e.printStackTrace();
        }
    }
}
@Component
public class MqttConsumerCallBack implements MqttCallback{
    /**
     * 客戶端斷開連接的回調(diào)
     */
    @Override
    public void connectionLost(Throwable throwable) {
        System.out.println("與服務(wù)器斷開連接,可重連");
    }
    /**
     * 消息到達(dá)的回調(diào)
     */
    @Override
    public void messageArrived(String topic, MqttMessage message) throws Exception {
        System.out.println(String.format("接收消息主題 : %s",topic));
        System.out.println(String.format("接收消息Qos : %d",message.getQos()));
        System.out.println(String.format("接收消息內(nèi)容 : %s",new String(message.getPayload())));
        if(topic.equals("data")){
  System.out.println(String.format("接收消息retained : %b",message.isRetained()));
        }
        System.out.println(String.format("接收消息retained : %b",message.isRetained()));
    }
    /**
     * 消息發(fā)布成功的回調(diào)
     */
    @Override
    public void deliveryComplete(IMqttDeliveryToken iMqttDeliveryToken) {
    }
}

接下來是發(fā)布消息部分:

@Configuration
@Slf4j
public class MqttProviderConfig {
    @Value("${spring.mqtt.username}")
    private String username;
    @Value("${spring.mqtt.password}")
    private String password;
    @Value("${spring.mqtt.url}")
    private String hostUrl;
    private String clientId = "provider_id";
    @Value("${spring.mqtt.default.topic}")
    private String defaultTopic;
    /**
     * 客戶端對(duì)象
     */
    private MqttClient client;
    /**
     * 在bean初始化后連接到服務(wù)器
     */
    @PostConstruct
    public void init(){
        connect();
    }
    /**
     * 客戶端連接服務(wù)端
     */
    public void connect(){
        try{
           // System.out.println("Connecting to MQTT server: " + hostUrl + " with client ID: " + clientId);
            //創(chuàng)建MQTT客戶端對(duì)象
            client = new MqttClient(hostUrl,clientId,new MemoryPersistence());
            //連接設(shè)置
            MqttConnectOptions options = new MqttConnectOptions();
            options.setCleanSession(true);
            //設(shè)置連接用戶名
            options.setUserName(username);
            //設(shè)置連接密碼
            options.setPassword(password.toCharArray());
            //設(shè)置超時(shí)時(shí)間,單位為秒
            options.setConnectionTimeout(100);
            //設(shè)置心跳時(shí)間 單位為秒,表示服務(wù)器每隔 1.5*20秒的時(shí)間向客戶端發(fā)送心跳判斷客戶端是否在線
            options.setKeepAliveInterval(20);
            //設(shè)置遺囑消息的話題,若客戶端和服務(wù)器之間的連接意外斷開,服務(wù)器將發(fā)布客戶端的遺囑信息
            options.setWill("willTopic",(clientId + "與服務(wù)器斷開連接").getBytes(),0,false);
            //設(shè)置回調(diào)
            client.setCallback(new MqttProviderCallBack());
            client.connect(options);
        } catch(MqttException e){
            e.printStackTrace();
        }
    }
    public void publish(int qos,boolean retained,String topic,String message){
        MqttMessage mqttMessage = new MqttMessage(); //創(chuàng)建消息實(shí)例
        mqttMessage.setQos(qos);  //設(shè)置qos
        mqttMessage.setRetained(retained);  //設(shè)置是否保留信息
        mqttMessage.setPayload(message.getBytes());
        //主題的目的地,用于發(fā)布/訂閱信息
        MqttTopic mqttTopic = client.getTopic(topic);
        //提供一種機(jī)制來跟蹤消息的傳遞進(jìn)度
        //用于在以非阻塞方式(在后臺(tái)運(yùn)行)執(zhí)行發(fā)布是跟蹤消息的傳遞進(jìn)度
        MqttDeliveryToken token;
        try {
            //將指定消息發(fā)布到主題,但不等待消息傳遞完成,返回的token可用于跟蹤消息的傳遞狀態(tài)
            token = mqttTopic.publish(mqttMessage);
            token.waitForCompletion();
        } catch (MqttException e) {
            e.printStackTrace();
        }
    }
}
@Configuration
public class MqttProviderCallBack implements MqttCallback{
    @Value("${spring.mqtt.client.id}")
    private String clientId;
    /**
     * 與服務(wù)器斷開的回調(diào)
     */
    @Override
    public void connectionLost(Throwable cause) {
        System.out.println(clientId+"與服務(wù)器斷開連接");
    }
    /**
     * 消息到達(dá)的回調(diào)
     */
    @Override
    public void messageArrived(String topic, MqttMessage message) throws Exception {
    }
    /**
     * 消息發(fā)布成功的回調(diào)
     */
    @Override
    public void deliveryComplete(IMqttDeliveryToken token) {
        IMqttAsyncClient client = token.getClient();
        System.out.println(client.getClientId()+"發(fā)布消息成功!");
    }
}

接下來就是調(diào)用publish方法,傳入相應(yīng)參數(shù),即可完成消息的發(fā)布。

這里有幾個(gè)參數(shù)要注意一下:

一、MQTT Topic(消息路由核心)

Topic 是客戶端發(fā)布 / 訂閱消息的 “地址標(biāo)識(shí)”,是實(shí)現(xiàn) “發(fā)布 - 訂閱” 模式的基礎(chǔ)。

  • 核心屬性:字符串格式,無(wú)預(yù)定義結(jié)構(gòu),由客戶端自定義,比如 “device/light/ 客廳”“data/sensor/ 溫度”。
  • 層級(jí)與通配符:用斜杠 “/” 劃分層級(jí),支持兩種通配符訂閱 ——“+” 匹配單個(gè)層級(jí)(如 “device/+/ 客廳” 匹配 “device/light/ 客廳”),“#” 匹配當(dāng)前及所有子層級(jí)(如 “data/#” 匹配所有數(shù)據(jù)類主題)。
  • 核心規(guī)則:發(fā)布者僅需指定 Topic 發(fā)送消息,訂閱者通過匹配 Topic(或通配符)接收消息,發(fā)布者與訂閱者無(wú)直接關(guān)聯(lián),實(shí)現(xiàn)解耦。

二、MQTT QoS(消息傳輸質(zhì)量等級(jí))

QoS(Quality of Service)定義了消息從發(fā)布者到訂閱者的傳輸可靠性,MQTT 3.1.1 標(biāo)準(zhǔn)規(guī)定了 3 個(gè)等級(jí),優(yōu)先級(jí)從低到高。

  1. QoS 0(最多一次):消息僅發(fā)送一次,不確認(rèn)、不重發(fā),可能丟失。適用于對(duì)可靠性要求低的場(chǎng)景,如實(shí)時(shí)溫度上報(bào)。
  2. QoS 1(至少一次):消息確保送達(dá),但可能重復(fù)。發(fā)布者發(fā)送后等待確認(rèn),未收到確認(rèn)則重發(fā),直到訂閱者確認(rèn)接收。
  3. QoS 2(恰好一次):消息確保僅送達(dá)一次,無(wú)丟失、無(wú)重復(fù)。通過 “發(fā)布 - 確認(rèn) - 釋放 - 完成” 四次握手實(shí)現(xiàn),適用于金融交易、指令下發(fā)等關(guān)鍵場(chǎng)景。

三、MQTT Retained

Retained 消息是 Broker(服務(wù)器)為指定 Topic 保存的 “最新一條消息”,具備 “狀態(tài)快照” 屬性。

  • 發(fā)布時(shí)觸發(fā):客戶端發(fā)布消息時(shí),需顯式設(shè)置 “Retain 標(biāo)志位” 為 true,Broker 才會(huì)保存該消息。
  • 僅存最新:同一 Topic 后續(xù)發(fā)布的 Retained 消息會(huì)覆蓋舊消息,Broker 始終只保留該 Topic 的最新狀態(tài)。
  • 在有別的客戶端訂閱該Topic時(shí),如果其“Retain 標(biāo)志位” 為 true,一旦連接,客戶端馬上會(huì)收到一條保存的最后發(fā)布的該Topic的消息。

總結(jié)

本文詳細(xì)描述了mqtt項(xiàng)目大致結(jié)構(gòu)、Sringboot整合mqtt代碼、mqtt的幾個(gè)重要參數(shù),希望對(duì)未來的架構(gòu)師們有幫助,謝謝~

到此這篇關(guān)于Springboot整合mqtt實(shí)現(xiàn)軟硬件通信的示例代碼的文章就介紹到這了,更多相關(guān)Springboot整合mqtt通信內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

最新評(píng)論

斗六市| 汝州市| 汾阳市| 木兰县| 防城港市| 河曲县| 康马县| 全南县| 海宁市| 平南县| 京山县| 莎车县| 上饶县| 延边| 海阳市| 萨迦县| 二连浩特市| 道孚县| 克什克腾旗| 梓潼县| 义乌市| 延长县| 白沙| 南澳县| 萨嘎县| 华亭县| 玉环县| 岑巩县| 同心县| 沙田区| 内黄县| 尉氏县| 武定县| 城固县| 定日县| 志丹县| 甘洛县| 简阳市| 富宁县| 宾阳县| 墨竹工卡县|