首页 > 编程语言 >SpringBoot整合MQTT实现低功耗物联网通信教程

SpringBoot整合MQTT实现低功耗物联网通信教程

来源:互联网 2026-07-26 07:56:03

MQTT是一种轻量级发布/订阅协议,专为低带宽、不可靠网络设计,具备低功耗、高并发优势。Mosquitto是其完整实现。SpringBoot整合Mosquitto需引入依赖、配置连接参数,并通过MqttConfig类管理客户端、主题及发布订阅通道,实现低功耗物联网通信。

mosquitto的简介

MQTT(Message Queuing Telemetry Transport,消息队列遥测传输协议)是一种轻量级的发布/订阅协议,专为网络条件苛刻的环境设计,能够在低带宽、不可靠或间歇性通信的场景下稳定运行。目前,MQTT已成为物联网消息通信领域的事实标准协议。

SpringBoot整合MQTT实现低功耗物联网通信教程

长期稳定更新的攒劲资源: >>>点此立即查看<<<

MQTT的突出特点之一是其提供的三种不同质量的消息服务。在工业物联网场景中,选择MQTT进行数据传输主要基于以下核心优势:

  1. 低协议开销:每条消息的消息头可短至2字节,网络开销极低,同时具备良好的容错性。
  2. 容错恢复能力强:物联网网络环境通常较为恶劣,MQTT能够从断开故障中自动恢复,无需额外编码。若使用HTTP,则需要自行实现重试逻辑。
  3. 低功耗:MQTT自设计之初便面向低功耗场景,非常适合电池供电的设备。
  4. 高并发支持:可接纳多达百万级别的客户端连接。

Mosquitto是MQTT协议的一个完整实现,也是业内广泛使用的选择。

近期工作中恰好使用SpringBoot整合Mosquitto,现将完整的搭建过程整理如下。项目基于Maven构建,以下是具体操作步骤。

1. 引入jar包依赖

首先,在pom.xml中添加以下依赖:


    org.springframework.integration
    spring-integration-core


    org.springframework.boot
    spring-boot-starter-integration


    org.springframework.integration
    spring-integration-stream


    org.springframework.integration
    spring-integration-mqtt

2. application.yml 配置文件

接下来编写配置文件,将MQTT的连接信息、主题、超时等参数放入其中:

spring:
  #mqtt配置
  mqtt:
    send:
      #完成超时时间
      completionTimeout: 3000
      #通过mqtt发送消息验证所需用户名
      username: test1
      #通过mqtt发送消息验证所需密码
      password: test1
      #连接的mqtt地址
      url: tcp: localhost:1883
      #客户端id
      clientId: clint1
      #推送主题  后面跟着#是监控下面所有的话题
      topic: topic
      #topic: my-test
      keepAliveInterval: 20
      connectionTimeout: 3000

3. MQTT配置文件(初始客户端、主题、发布、订阅等)

创建MqttConfig类,集中管理MQTT的所有配置,包括客户端初始化、主题绑定、发布和订阅通道等:

import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
import org.eclipse.paho.client.mqttv3.MqttMessage;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.integration.annotation.ServiceActivator;
import org.springframework.integration.channel.DirectChannel;
import org.springframework.integration.core.MessageProducer;
import org.springframework.integration.mqtt.core.DefaultMqttPahoClientFactory;
import org.springframework.integration.mqtt.core.MqttPahoClientFactory;
import org.springframework.integration.mqtt.inbound.MqttPahoMessageDrivenChannelAdapter;
import org.springframework.integration.mqtt.outbound.MqttPahoMessageHandler;
import org.springframework.integration.mqtt.support.DefaultPahoMessageConverter;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.MessageHandler;
import org.springframework.util.StringUtils;

/**
 * @program: mg_parse
 * @description: TODO  MQTT配置,生产者
 * @author: 
 * @create: 
 **/
@Configuration
public class MqttConfig {
    private static final byte[] WILL_DATA;
    static {
        WILL_DATA = "offline".getBytes();
    }
    /**
     * mqtt订阅者使用信道名称
     */
    public static final String CHANNEL_NAME_IN = "mqttInboundChannel";
    /**
     * mqtt发布者信道名称
     */
    public static final String CHANNEL_NAME_OUT = "mqttOutboundChannel";
    /**
     * mqtt发送者用户名
     */
    @Value("${spring.mqtt.send.username}")
    private String username;
    /**
     * mqtt发送者密码
     */
    @Value("${spring.mqtt.send.password}")
    private String password;
    /**
     * mqtt发送者url
     */
    @Value("${spring.mqtt.send.url}")
    private String hostUrl;
    /**
     * mqtt发送者客户端id
     */
    @Value("${spring.mqtt.send.clientId}")
    private String clientId;
    /**
     * mqtt发送者主题
     */
    @Value("${spring.mqtt.send.topic}")
    private String msgTopic;
    /**
     * mqtt发送者超时时间
     */
    @Value("${spring.mqtt.send.completionTimeout}")
    private int completionTimeout ;
    /**
     * @author liujianfu
     * @description
     */
    @Value("${spring.mqtt.send.keepAliveInterval}")
    private int keepAliveInterval;
    /**
     * @author liujianfu
     * @description
     */
    @Value("${spring.mqtt.send.connectionTimeout}")
    private int connectionTimeout;

    @Autowired
    private MqttCallbackHandler mqttCallbackHandler;


    /**
     * @author liujianfu
     * @description   新建MqttConnectionOptionsBean  MQTT连接器选项
     * @date 2021/8/17 10:34
     * @param
     * @return org.eclipse.paho.client.mqttv3.MqttConnectOptions
     */
    @Bean
    public MqttConnectOptions getSenderMqttConnectOptions(){
        MqttConnectOptions options=new MqttConnectOptions();
        // 设置连接的用户名
        if(!username.trim().equals("")){
            //将用户名去掉前后空格
            options.setUserName(username);
        }
        // 设置连接的密码
        options.setPassword(password.toCharArray());
        // 转化连接的url地址
        String[] uris={hostUrl};
        // 设置连接的地址
        options.setServerURIs(uris);
        // 设置超时时间 单位为秒
        options.setConnectionTimeout(completionTimeout);
        // 设置会话心跳时间 单位为秒 服务器会每隔1.5*20秒的时间向客户端发送心跳判断客户端是否在线
        // 但这个方法并没有重连的机制
        options.setKeepAliveInterval(keepAliveInterval);
        // 设置“遗嘱”消息的话题,若客户端与服务器之间的连接意外中断,服务器将发布客户端的“遗嘱”消息。
        //设置超时时间
        options.setConnectionTimeout(connectionTimeout);
        options.setCleanSession(true);
        options.setAutomaticReconnect(true);
        return options;
    }
    /**
     *创建MqttPathClientFactoryBean
     */
    @Bean
    public MqttPahoClientFactory senderMqttClientFactory() {
        //创建mqtt客户端工厂
        DefaultMqttPahoClientFactory factory = new DefaultMqttPahoClientFactory();
        //设置mqtt的连接设置
        factory.setConnectionOptions(getSenderMqttConnectOptions());
        return factory;
    }
    /**
     * 发布者-MQTT信息通道(生产者)
     */
    @Bean(name = CHANNEL_NAME_OUT)
    public MessageChannel mqttOutboundChannel() {
        return new DirectChannel();
    }


    /**
     * 发布者-MQTT消息处理器(生产者)  将channel绑定到MqttClientFactory上
     *
     * @return {@link org.springframework.messaging.MessageHandler}
     */
    @Bean
    @ServiceActivator(inputChannel = CHANNEL_NAME_OUT)
    public MessageHandler mqttOutbound() {
        //创建消息处理器
        MqttPahoMessageHandler messageHandler = new MqttPahoMessageHandler(
                clientId+"_pub",
                senderMqttClientFactory());
        //设置消息处理类型为异步
        messageHandler.setAsync(true);
        //设置消息的默认主题
        messageHandler.setDefaultTopic(msgTopic);
        messageHandler.setDefaultRetained(false);
        //1.重新连接MQTT服务时,不需要接收该主题最新消息,设置retained为false;
        //2.重新连接MQTT服务时,需要接收该主题最新消息,设置retained为true;
        return messageHandler;
    }



    /************                             消费者,订阅者的消费信息                               *****/

    /**
     * MQTT信息通道(消费者)
     *
     */
    @Bean(name = CHANNEL_NAME_IN)
    public MessageChannel mqttInboundChannel() {
        return new DirectChannel();
    }
    /**
     * MQTT消息订阅绑定(消费者)
     *
     */
    @Bean
    public MessageProducer inbound() {
        System.out.println("topics:"+msgTopic);
        // 可以同时消费(订阅)多个Topic
        MqttPahoMessageDrivenChannelAdapter adapter =
                new MqttPahoMessageDrivenChannelAdapter(
                        clientId+"_sub", senderMqttClientFactory(), msgTopic);
        adapter.setCompletionTimeout(5000);
        adapter.setConverter(new DefaultPahoMessageConverter());
        adapter.setQos(0);
        // 设置订阅通道
        adapter.setOutputChannel(mqttInboundChannel());
        return adapter;
    }

    /**
     * MQTT消息处理器(消费者)
     *
     */
    @Bean
    @ServiceActivator(inputChannel = CHANNEL_NAME_IN)
    public MessageHandler handler() {
        return message -> {
            String topic = message.getHeaders().get("mqtt_receivedTopic").toString();
            String payload = message.getPayload().toString();
            mqttCallbackHandler.handle(topic,payload);
        };
    }
}

4. 发送接口网关MqSendMessageGateWay

创建一个消息发送网关,通过Spring Integration的@MessagingGateway注解,将消息发送到指定的MQTT通道:

import org.springframework.integration.annotation.MessagingGateway;
import org.springframework.integration.mqtt.support.MqttHeaders;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.stereotype.Component;

/**
 * @program: mg_parse
 * @description:
 * @author: 
 * @create:
 **/

@Component
@MessagingGateway(defaultRequestChannel = MqttConfig.CHANNEL_NAME_OUT)
public interface MqSendMessageGateWay {
    /**
     * 默认的消息机制
     * @param data
     */
    void sendToMqtt(String data);

    /**
     * 发送消息 向mqtt指定topic发送消息
     * @param topic
     * @param payload
     */
    void sendToMqtt(@Header(MqttHeaders.TOPIC) String topic, String payload);

    /**
     * 发送消息 向mqtt指定topic发送消息
     * @param topic
     * @param qos
     * @param payload
     */
    void sendToMqtt(@Header(MqttHeaders.TOPIC) String topic, @Header(MqttHeaders.QOS) int qos, String payload);
}

5. 消息订阅类MqttCallbackHandler

订阅消息后,需要一个回调处理器来接收并处理消息。这里定义了一个简单的处理器,根据topic对消息进行分发处理:

import org.springframework.stereotype.Service;

/**
 * @program: mg_parse
 * @description:
 * @author: 
 * @create: 
 **/
@Service
public class MqttCallbackHandler {

    public void handle(String topic, String payload){
        // 根据topic分别进行消息处理。
        System.out.println("MqttCallbackHandle:" + topic + "|"+ payload);
    }
}

6. 发布生产,这里用controller层来测试

在实际生产中,通常是在Service层中调用网关发送消息。这里使用一个Controller来演示测试:

@Autowired
private MqSendMessageGateWay mqSendMessageGateWay;

@RequestMapping("/send")
@ResponseBody
private ResponseEntity send(){
    String data = "我是springboot发送的数据";
    mqSendMessageGateWay.sendToMqtt(data);
    // return R.ok("OK");
    return new ResponseEntity<>("OK", HttpStatus.OK);
}
/**
 * 动态增加主题
 * @param
 * @param
 */
@ResponseBody
@RequestMapping("/sendToTopic")
private ResponseEntity sendToTopic(){
    String topic = "sharjeck/ai/test/out";
    String data = "这是出的主题";
    mqSendMessageGateWay.sendToMqtt(topic,data);
    return new ResponseEntity<>("OK", HttpStatus.OK);
}

7. 启动SpringBoot启动类即可

8. SpringBoot连接MQTT进行发布消息时取消保留历史消息

在发布消息后,如果设置不当,新的订阅者可能会收到之前发布的历史消息。如果业务上只需要接收即时消息,而不需要历史数据,则需要关注一个参数——retained。

Retained 消息是指在PUBLISH数据包中Retain标识设为1的消息。Broker收到此类消息后,会将其保存下来,当有新的订阅者订阅相应主题时,Broker会立即将该消息推送给订阅者。

因此,在生产者端,只需将retained设置为false,即可避免保留历史消息。至于QoS等级,可根据实际业务需求灵活设定。至此,配置基本完成。

总结

以上就是完整的SpringBoot整合MQTT的实战配置。从协议简介到具体代码实现,覆盖了生产者和消费者的核心环节,希望能为大家提供一个实用的参考。

侠游戏发布此文仅为了传递信息,不代表侠游戏网站认同其观点或证实其描述

热游推荐

更多
湘ICP备14008430号-1 湘公网安备 43070302000280号
All Rights Reserved
本站为非盈利网站,不接受任何广告。本站所有软件,都由网友
上传,如有侵犯你的版权,请发邮件给xiayx666@163.com
抵制不良色情、反动、暴力游戏。注意自我保护,谨防受骗上当。
适度游戏益脑,沉迷游戏伤身。合理安排时间,享受健康生活。