首页技术栈归档照片墙音乐日记随想收藏夹友链留言关于

MQTT 消息协议详解

写作时间:2026-07-07

一、MQTT简介

MQTT(Message Queuing Telemetry Transport,消息队列遥测传输协议),是一种基于发布/订阅(publish/subscribe)模式的"轻量级"通讯协议,该协议构建于TCP/IP协议上,由IBM在1999年发布。MQTT最大优点在于,可以以极少的代码和有限的带宽,为连接远程设备提供实时可靠的消息服务。作为一种低开销、低带宽占用的即时通讯协议,使其在物联网、小型设备、移动应用等方面有较广泛的应用。

MQTT是一个基于客户端-服务器的消息发布/订阅传输协议。MQTT协议是轻量、简单、开放和易于实现的,这些特点使它适用范围非常广泛。在很多情况下,包括受限的环境中,如:机器与机器(M2M)通信和物联网(IoT)。其在,通过卫星链路通信传感器、偶尔拨号的医疗设备、智能家居、及一些小型化设备中已广泛使用。

二、设计规范

由于物联网的环境是非常特别的,所以MQTT遵循以下设计原则:

  • (1)精简,不添加可有可无的功能;
  • (2)发布/订阅(Pub/Sub)模式,方便消息在传感器之间传递;
  • (3)允许用户动态创建主题,零运维成本;
  • (4)把传输量降到最低以提高传输效率;
  • (5)把低带宽、高延迟、不稳定的网络等因素考虑在内;
  • (6)支持连续的会话控制;
  • (7)理解客户端计算能力可能很低;
  • (8)提供服务质量管理;
  • (9)假设数据不可知,不强求传输数据的类型与格式,保持灵活性。

三、主要特性

MQTT协议工作在低带宽、不可靠的网络的远程传感器和控制设备通讯而设计的协议,它具有以下主要的几项特性:

  • (1)使用发布/订阅消息模式,提供一对多的消息发布,解除应用程序耦合。

    这一点很类似于XMPP,但是MQTT的信息冗余远小于XMPP,,因为XMPP使用XML格式文本来传递数据。

  • (2)对负载内容屏蔽的消息传输。

  • (3)使用TCP/IP提供网络连接。

    主流的MQTT是基于TCP连接进行数据推送的,但是同样有基于UDP的版本,叫做MQTT-SN。这两种版本由于基于不同的连接方式,优缺点自然也就各有不同了。

  • (4)有三种消息发布服务质量:

    "至多一次",消息发布完全依赖底层TCP/IP网络。会发生消息丢失或重复。这一级别可用于如下情况,环境传感器数据,丢失一次读记录无所谓,因为不久后还会有第二次发送。这一种方式主要普通APP的推送,倘若你的智能设备在消息推送时未联网,推送过去没收到,再次联网也就收不到了。

    "至少一次",确保消息到达,但消息重复可能会发生。

    "只有一次",确保消息到达一次。在一些要求比较严格的计费系统中,可以使用此级别。在计费系统中,消息重复或丢失会导致不正确的结果。这种最高质量的消息发布服务还可以用于即时通讯类的APP的推送,确保用户收到且只会收到一次。

  • (5)小型传输,开销很小(固定长度的头部是2字节),协议交换最小化,以降低网络流量。

    这就是为什么在介绍里说它非常适合"在物联网领域,传感器与服务器的通信,信息的收集",要知道嵌入式设备的运算能力和带宽都相对薄弱,使用这种协议来传递消息再适合不过了。

  • (6)使用Last Will和Testament特性通知有关各方客户端异常中断的机制。

    Last Will:即遗言机制,用于通知同一主题下的其他设备发送遗言的设备已经断开了连接。

    Testament:遗嘱机制,功能类似于Last Will。

四、MQTT协议原理

4.1 MQTT协议实现方式

实现MQTT协议需要客户端和服务器端通讯完成,在通讯过程中,MQTT协议中有三种身份:发布者(Publish)、代理(Broker)(服务器)、订阅者(Subscribe)。其中,消息的发布者和订阅者都是客户端,消息代理是服务器,消息发布者可以同时是订阅者。

MQTT传输的消息分为:主题(Topic)和负载(payload)两部分:

  • (1)Topic,可以理解为消息的类型,订阅者订阅(Subscribe)后,就会收到该主题的消息内容(payload);
  • (2)payload,可以理解为消息的内容,是指订阅者具体要使用的内容。

4.2 网络传输与应用消息

MQTT会构建底层网络传输:它将建立客户端到服务器的连接,提供两者之间的一个有序的、无损的、基于字节流的双向传输。

当应用数据通过MQTT网络发送时,MQTT会把与之相关的服务质量(QoS)和主题名(Topic)相关连。

4.3 MQTT客户端

一个使用MQTT协议的应用程序或者设备,它总是建立到服务器的网络连接。客户端可以:

  • (1)发布其他客户端可能会订阅的信息;
  • (2)订阅其它客户端发布的消息;
  • (3)退订或删除应用程序的消息;
  • (4)断开与服务器连接。

4.4 MQTT服务器

MQTT服务器以称为"消息代理"(Broker),可以是一个应用程序或一台设备。它是位于消息发布者和订阅者之间,它可以:

  • (1)接受来自客户的网络连接;
  • (2)接受客户发布的应用信息;
  • (3)处理来自客户端的订阅和退订请求;
  • (4)向订阅的客户转发应用程序消息。

4.5 MQTT协议中的订阅、主题、会话

一、订阅(Subscription)

订阅包含主题筛选器(Topic Filter)和最大服务质量(QoS)。订阅会与一个会话(Session)关联。一个会话可以包含多个订阅。每一个会话中的每个订阅都有一个不同的主题筛选器。

二、会话(Session)

每个客户端与服务器建立连接后就是一个会话,客户端和服务器之间有状态交互。会话存在于一个网络之间,也可能在客户端和服务器之间跨越多个连续的网络连接。

三、主题名(Topic Name)

连接到一个应用程序消息的标签,该标签与服务器的订阅相匹配。服务器会将消息发送给订阅所匹配标签的每个客户端。

四、主题筛选器(Topic Filter)

一个对主题名通配符筛选器,在订阅表达式中使用,表示订阅所匹配到的多个主题。

五、负载(Payload)

消息订阅者所具体接收的内容。

4.6 MQTT协议中的方法

MQTT协议中定义了一些方法(也被称为动作),来于表示对确定资源所进行操作。这个资源可以代表预先存在的数据或动态生成数据,这取决于服务器的实现。通常来说,资源指服务器上的文件或输出。主要方法有:

  • (1)Connect。等待与服务器建立连接。
  • (2)Disconnect。等待MQTT客户端完成所做的工作,并与服务器断开TCP/IP会话。
  • (3)Subscribe。等待完成订阅。
  • (4)UnSubscribe。等待服务器取消客户端的一个或多个topics订阅。
  • (5)Publish。MQTT客户端发送消息请求,发送完成后返回应用程序线程。

五、MQTT协议数据包结构

在MQTT协议中,一个MQTT数据包由:固定头(Fixed header)、可变头(Variable header)、消息体(payload)三部分构成。MQTT数据包结构如下:

  • (1)固定头(Fixed header)。存在于所有MQTT数据包中,表示数据包类型及数据包的分组类标识。
  • (2)可变头(Variable header)。存在于部分MQTT数据包中,数据包类型决定了可变头是否存在及其具体内容。
  • (3)消息体(Payload)。存在于部分MQTT数据包中,表示客户端收到的具体内容。

5.1 MQTT固定头

固定头存在于所有MQTT数据包中,其结构如下:

5.1.1 MQTT数据包类型

位置:Byte 1中bits 7-4。

相于一个4位的无符号值,类型、取值及描述如下:

5.1.2 标识位

位置:Byte 1中bits 3-0。

在不使用标识位的消息类型中,标识位被作为保留位。如果收到无效的标志时,接收端必须关闭网络连接:

(1)DUP:发布消息的副本。用来在保证消息的可靠传输,如果设置为1,则在下面的变长中增加MessageId,并且需要回复确认,以保证消息传输完成,但不能用于检测消息重复发送。

(2)QoS:发布消息的服务质量,即:保证消息传递的次数

Ø00:最多一次,即:<=1

Ø01:至少一次,即:>=1

Ø10:一次,即:=1

Ø11:预留

(3)RETAIN: 发布保留标识,表示服务器要保留这次推送的信息,如果有新的订阅者出现,就把这消息推送给它,如果设有那么推送至当前订阅者后释放。

5.1.3 剩余长度(Remaining Length)

地址:Byte 2。

固定头的第二字节用来保存变长头部和消息体的总大小的,但不是直接保存的。这一字节是可以扩展,其保存机制,前7位用于保存长度,后一部用做标识。当最后一位为1时,表示长度不足,需要使用二个字节继续保存。例如:计算出后面的大小为0

5.2 MQTT可变头

MQTT数据包中包含一个可变头,它驻位于固定的头和负载之间。可变头的内容因数据包类型而不同,较常的应用是作为包的标识:

很多类型数据包中都包括一个2字节的数据包标识字段,这些类型的包有:PUBLISH (QoS > 0)、PUBACK、PUBREC、PUBREL、PUBCOMP、SUBSCRIBE、SUBACK、UNSUBSCRIBE、UNSUBACK。

5.3 Payload消息体

Payload消息体位MQTT数据包的第三部分,包含CONNECT、SUBSCRIBE、SUBACK、UNSUBSCRIBE四种类型的消息:

  • (1)CONNECT,消息体内容主要是:客户端的ClientID、订阅的Topic、Message以及用户名和密码。
  • (2)SUBSCRIBE,消息体内容是一系列的要订阅的主题以及QoS。
  • (3)SUBACK,消息体内容是服务器对于SUBSCRIBE所申请的主题及QoS进行确认和回复。
  • (4)UNSUBSCRIBE,消息体内容是要订阅的主题

MQTT基本使用

直接解压在bin目录下执行命令

mqtt命令

#启动MQTT
emqx.cmd start
#停止MQTT
emqx.cmd stop

服务端运行后,可以在浏览器中输入地址http://127.0.0.1:18083 进入后台管理,用户名为admin,密码为public

登录后我们主要看这三个菜单

**1.Clients:**当前连接的客户端列表

**2.Topics:**订阅主题列表

**3.subscriptions:**订阅用户列表

接下来,你可以用本地客户端连接服务端来进行验证。

二、mqtt本地客户端安装 emq提供了在线web客户端,可以用来连接到emq提供的服务端进行验证和测试。

但在开发环境和生产环境,我们需要部署本地客户端,连接到我们本地服务器上进行调试。

1.连接服务端

运行客户端mqttx程序,点击添加new connection,录入连接名称和服务端IP(此处为连接本机服务端),其他选项不用改,点击connect后,客户端成功连接到服务端。

测试

导入依赖

		<dependency>
			<groupId>org.eclipse.paho</groupId>
			<artifactId>org.eclipse.paho.client.mqttv3</artifactId>
			<version>1.2.0</version>
		</dependency>
		<dependency>
			<groupId>org.projectlombok</groupId>
			<artifactId>lombok</artifactId>
		</dependency>

在单线程的环境下调用MQTT

package com.baosight.mgrzfn;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.paho.client.mqttv3.*;
import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;

/**
 * 在单线程的环境下使用
 */
@Slf4j
public class DemoMqttClient {


    public static void main(String... args) {
        try {
            // host为主机名,clientid即连接MQTT的客户端ID,一般以客户端唯一标识符表示,
            // MemoryPersistence设置clientid的保存形式,默认为以内存保存
            MqttClient mqttClient = new MqttClient("tcp://127.0.0.1:1883", "client", new MemoryPersistence());
            // 配置参数信息
            MqttConnectOptions options = new MqttConnectOptions();
            // 设置是否清空session,这里如果设置为false表示服务器会保留客户端的连接记录,
            // 这里设置为true表示每次连接到服务器都以新的身份连接
            options.setCleanSession(true);
            // 设置用户名
            options.setUserName("admin");
            // 设置密码
            options.setPassword("public".toCharArray());
            // 设置超时时间 单位为秒
            options.setConnectionTimeout(10);
            // 设置会话心跳时间 单位为秒 服务器会每隔1.5*20秒的时间向客户端发送个消息判断客户端是否在线,但这个方法并没有重连的机制
            options.setKeepAliveInterval(20);
            // 连接
            mqttClient.connect(options);
            // 订阅
            mqttClient.subscribe("test");
            // 设置回调
            mqttClient.setCallback(new MqttCallback() {
                @Override
                public void connectionLost(Throwable throwable) {
                    // 连接失败时调用  重新连接订阅
                    System.out.println("连接丢失.............");
                    try {
                        System.out.println("开始重连");
                        Thread.sleep(3000);
                        mqttClient.connect(options);
                    } catch (InterruptedException e) {
                        e.printStackTrace();
                    } catch (MqttSecurityException e) {
                        e.printStackTrace();
                    } catch (MqttException e) {
                        e.printStackTrace();
                    }
                }

                @Override
                public void messageArrived(String topic, MqttMessage mqttMessage) throws Exception {
                    log.info("接收消息主题 : " + topic);
                    log.info("接收消息Qos : " + mqttMessage.getQos());
                    log.info("接收消息内容 : " + new String(mqttMessage.getPayload()));
                }

                @Override
                public void deliveryComplete(IMqttDeliveryToken iMqttDeliveryToken) {
                    //认证过程
                    log.info("deliveryComplete.............");
                }
            });


            // 创建消息,给我MQTT进行消息的发送
            MqttMessage message = new MqttMessage("{\"id\":\"123\", \"name\":\"abcd\"}".getBytes());
            // 设置消息的服务质量
            message.setQos(0);
            // 发布消息
            mqttClient.publish("test", message);
            // 断开连接
            /**
             * 可以注释掉停止连接来查看一下三个是否有数据,如果关闭了连接就查询不到,没关闭可以查询到
             * 1.Clients:当前连接的客户端列表
             * 2.Topics:订阅主题列
             * 3.subscriptions:订阅用户列表
             */
            mqttClient.disconnect();
            // 关闭客户端
            mqttClient.close();
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

在多线程的环境下调用MQTT

package com.baosight.mgrzfn;


import lombok.extern.slf4j.Slf4j;
import org.eclipse.paho.client.mqttv3.*;
import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;

/**
 *在多线程的环境下运行MQTT
 */
@Slf4j
public class MqttClientThread extends Thread{

    //连接地址
    private String serverURL;
    //MQTT客户端登录用户名
    private String mqttUsername;
    //MQTT客户端密码
    private String mqttPassWord;
    //MQTT订阅主题
    private String mqttTopic;
    //MQTT的client
    private String clientId;
    //产品id
    private String productId;
    //推送至我们自己的RedisTopIc中channel
    private String channel = "mqtt";
    //mqtt实体类
    private MqttClient mqttClient;

    //构造函数
    public MqttClientThread(String serverURL,String mqttUsername,String mqttPassWord,String mqttTopic,String clientId,String productId) {
        this.serverURL = serverURL;
        this.mqttUsername = mqttUsername;
        this.mqttPassWord = mqttPassWord;
        this.mqttTopic = mqttTopic;
        this.clientId = clientId;
        this.productId = productId;
    }

    //线程方法
    public void run(){
        try {
            // host为主机名,clientid即连接MQTT的客户端ID,一般以客户端唯一标识符表示,
            // MemoryPersistence设置clientid的保存形式,默认为以内存保存,就用username
            mqttClient = new MqttClient(serverURL, clientId, new MemoryPersistence());
            // 配置参数信息
            MqttConnectOptions options = new MqttConnectOptions();
            // 设置是否清空session,这里如果设置为false表示服务器会保留客户端的连接记录,
            // 这里设置为true表示每次连接到服务器都以新的身份连接
            options.setCleanSession(true);
            // 设置用户名
            options.setUserName(mqttUsername);
            // 设置密码
            options.setPassword(mqttPassWord.toCharArray());
            // 设置超时时间 单位为秒
            options.setConnectionTimeout(10);
            // 设置会话心跳时间 单位为秒 服务器会每隔1.5*20秒的时间向客户端发送个消息判断客户端是否在线,但这个方法并没有重连的机制
//            options.setKeepAliveInterval(20);
            //设置断开后重新连接
            options.setAutomaticReconnect(true);
            // 连接
            mqttClient.connect(options);
            // 订阅
            //如果监测到有,号,说明要订阅多个主题
            if(mqttTopic.contains(",")){
                //多主题
                String[] mqttTopics = mqttTopic.split(",");
                mqttClient.subscribe(mqttTopics);
            }else{
                //单主题
                mqttClient.subscribe(mqttTopic);
            }
            // 设置回调
            mqttClient.setCallback(new MqttCallbackExtended () {
                /**
                 * Called when the connection to the server is completed successfully.
                 *
                 * @param reconnect If true, the connection was the result of automatic reconnect.
                 * @param serverURI The server URI that the connection was made to.
                 */
                @Override
                public void connectComplete(boolean reconnect, String serverURI) {
                    try{
                        //如果监测到有,号,说明要订阅多个主题
                        if(mqttTopic.contains(",")){
                            //多主题
                            String[] mqttTopics = mqttTopic.split(",");
                            mqttClient.subscribe(mqttTopics);
                        }else{
                            //单主题
                            mqttClient.subscribe(mqttTopic);
                        }
                        log.info("----TAG", "connectComplete: 订阅主题成功");
                    }catch(Exception e){
                        e.printStackTrace();
                        log.info("----TAG", "error: 订阅主题失败");
                    }
                }


                @Override
                public void connectionLost(Throwable throwable) {
                    log.error("连接断开,下面做重连...");
                    long reconnectTimes = 1;
                    while (true) {
                        try {
                            if (mqttClient.isConnected()) {
                                log.warn("mqtt reconnect success end");
                                break;
                            }
                            if(reconnectTimes == 10){
                                //当重连次数达到10次时,就抛出异常,不在重连
                                log.warn("mqtt reconnect error");
                                return;
                            }
                            log.warn("mqtt reconnect times = {} try again...", reconnectTimes++);
                            mqttClient.reconnect();
                        } catch (MqttException e) {
                            log.error("", e);
                        }
                        try {
                            Thread.sleep(1000);
                        } catch (InterruptedException e1) {
//                            e1.printStackTrace();
                        }
                    }
                }

                @Override
                public void messageArrived(String topic, MqttMessage mqttMessage) throws Exception {
                    log.info("接收消息主题 : " + topic);
                    log.info("接收消息Qos : " + mqttMessage.getQos());
                    log.info("接收消息内容 : " + new String(mqttMessage.getPayload()));
                    //向我们通道中发送消息
                    RedisMsg redisMsg = new RedisMsg();
                    redisMsg.setChannel(channel);
                    redisMsg.setMsg("推送MQTT消息");
                    SensorMqttMsg mqttMsg = new SensorMqttMsg();
                    mqttMsg.setProductId(productId);
                    mqttMsg.setPayload(new String(mqttMessage.getPayload()));
                    redisMsg.setData(mqttMsg);
                    RedisTopicUtil.sendMessage(channel, redisMsg);
                }

                @Override
                public void deliveryComplete(IMqttDeliveryToken iMqttDeliveryToken) {
                    //认证过程
                    log.info("deliveryComplete.............");
                }
            });
            //放入缓存,根据clinetId吧mqttClient对象放进去
//            MqttClientManager.MQTT_CLIENT_MAP.putIfAbsent(clientId, mqttClient);
        } catch (Exception e) {
            e.printStackTrace();
            //当创建客户端的时候出现 已断开连接,有可能是在另一个环境下启动了该客户端,直接吧这边的客户端关闭,不然另一边会无限重连
            if(e.getMessage().equals("已断开连接") || e.getMessage().equals("客户机未连接")){
                try {
                    mqttClient.close();
                } catch (MqttException ex) {
                    ex.printStackTrace();
                }
            }
        }
    }
}

MQTT工具类实现

导入依赖

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

application.properties

#MQTT配置信息
#MQTT-用户名
spring.mqtt.username=admin
#MQTT-密码
spring.mqtt.password=public
#MQTT-服务器连接地址,如果有多个,用逗号隔开,如:tcp://127.0.0.1:61613,tcp://47.123.33.66:61613
spring.mqtt.url=tcp://127.0.0.1:1883
#MQTT-连接服务器默认客户端ID,名字自己取,不重复即可
spring.mqtt.client.id=client
#MQTT-默认的消息推送主题,实际可在调用接口时指定
spring.mqtt.default.topic=test
#timeout 链接超时时间
mqtt.connection.timeout=20
#keep alive
mqtt.keep.alive.interval=20

config

package com.example.demo.config;

import lombok.extern.slf4j.Slf4j;
import org.eclipse.paho.client.mqttv3.IMqttDeliveryToken;
import org.eclipse.paho.client.mqttv3.MqttCallback;
import org.eclipse.paho.client.mqttv3.MqttMessage;
import org.springframework.stereotype.Component;

@Component
@Slf4j
public class MqttMessageCallback implements MqttCallback {
    /**
     * 链接丢失时处理
     * @param throwable
     */
    @Override
    public void connectionLost(Throwable throwable) {
        //可以做重连 或者 其他业务处理
    }

    @Override
    public void messageArrived(String topic, MqttMessage mqttMessage) throws Exception {
        log.info("接收到消息topic---->{}",topic);
        log.info("接收到消息质量qos---->{}",mqttMessage.getQos());
        log.info("接收到消息具体信息---->{}",new String(mqttMessage.getPayload()));

        System.out.println("订阅MQTT消息");
        System.out.println(new String(mqttMessage.getPayload()));
        //结合业务 编写具体信息即可
    }

    @Override
    public void deliveryComplete(IMqttDeliveryToken iMqttDeliveryToken) {

    }
}

utlis

package com.example.demo.utlis;

import com.example.demo.config.MqttMessageCallback;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.paho.client.mqttv3.*;
import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;

import javax.annotation.PostConstruct;
import java.util.Objects;

@Component
@Slf4j
public class MqttClientUtil {
    @Value("${spring.mqtt.username}")
    private String username;
    @Value("${spring.mqtt.password}")
    private String password;
    @Value("${spring.mqtt.url}")
    private String host;
    @Value("${spring.mqtt.client.id}")
    private String clientId;
    @Value("${spring.mqtt.default.topic}")
    private String topic;
    @Value("${mqtt.connection.timeout}")
    private int timeOut;
    @Value("${mqtt.keep.alive.interval}")
    private int interval;

    @Autowired
    private MqttMessageCallback mqttMessageCallback;

    private MqttClient mqttClient;
    private MqttConnectOptions mqttConnectOptions;

    @PostConstruct
    private void init(){
        connect(host, clientId,topic);
    }

    /**
     * 链接mqtt
     * @param host
     * @param clientId
     */
    private void connect(String host,String clientId,String topic){
        try{
            mqttClient = new MqttClient(host,clientId,new MemoryPersistence());
            mqttConnectOptions = getMqttConnectOptions();
            //设置回调函数
            mqttClient.setCallback(mqttMessageCallback);
            //链接mqtt
            mqttClient.connect(mqttConnectOptions);
            //订阅消息
            mqttClient.subscribe(topic,1);
        }catch (Exception e){
            log.error("mqtt服务链接异常!");
            e.printStackTrace();
        }
    }

    /**
     * 设置链接对象信息
     * setCleanSession  true 断开链接即清楚会话  false 保留链接信息 离线还会继续发消息
     * @return
     */
    private MqttConnectOptions getMqttConnectOptions(){
        MqttConnectOptions mqttConnectOptions = new MqttConnectOptions();
        mqttConnectOptions.setUserName(username);
        mqttConnectOptions.setPassword(password.toCharArray());
        mqttConnectOptions.setServerURIs(new String[]{host});
        mqttConnectOptions.setKeepAliveInterval(interval);
        mqttConnectOptions.setConnectionTimeout(timeOut);
        mqttConnectOptions.setCleanSession(true);
        return mqttConnectOptions;
    }

    /**
     *mqtt链接状态
     * @return
     */
    private boolean isConnect(){
        if(Objects.isNull(this.mqttClient)){
            return false;
        }
        return mqttClient.isConnected();
    }

    /**
     * 设置重连
     * @throws Exception
     */
    private void reConnect() throws Exception{
        if(Objects.nonNull(this.mqttClient)){
            log.info("mqtt 服务已重新链接...");
            this.mqttClient.connect(this.mqttConnectOptions);
        }
    }

    /**
     * 断开链接
     * @throws Exception
     */
    private void closeConnect() throws Exception{
        if(Objects.nonNull(this.mqttClient)){
            log.info("mqtt 服务已断开链接...");
            this.mqttClient.disconnect();
        }
    }

    /**
     * 发布消息
     * @param topic
     * @param message
     * @param qos
     * @throws Exception
     */
    public void sendMessage(String topic,String message,int qos) throws Exception {
        if(Objects.nonNull(this.mqttClient) && this.mqttClient.isConnected()){
            MqttMessage mqttMessage = new MqttMessage();
            mqttMessage.setPayload(message.getBytes());
            mqttMessage.setQos(qos);

            MqttTopic mqttTopic = mqttClient.getTopic(topic);

            if(Objects.nonNull(mqttTopic)){
                try{
                    MqttDeliveryToken publish = mqttTopic.publish(mqttMessage);
                    if(publish.isComplete()){
                        log.info("消息发送成功---->{}",message);
                    }
                }catch(Exception e){
                    log.error("消息发送异常",e);
                }
            }
        }else{
            reConnect();
        }
    }

}

controller

package com.example.demo.controller;


import com.example.demo.utlis.MqttClientUtil;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;

@RestController
@RequestMapping("/mqtt")
public class mqttController {
    @Autowired
    private MqttClientUtil mqttClientUtil;

    @GetMapping("/message/{msg}")
    public String sendMsg(@PathVariable("msg") String msg) throws Exception{
        mqttClientUtil.sendMessage("test",msg,1);
        return "OK";
    }
}
avatar

yuanyourdomain

写代码,做研究,记录生活。

RECOMMENDED

MyBatis 动态 SQL

2026-07-08

Nginx 基础入门

2026-07-08

Maven 多模块与私服

2026-07-08