MQTT是一种轻量级发布/订阅协议,专为低带宽、不可靠网络设计,具备低功耗、高并发优势。Mosquitto是其完整实现。SpringBoot整合Mosquitto需引入依赖、配置连接参数,并通过MqttConfig类管理客户端、主题及发布订阅通道,实现低功耗物联网通信。
MQTT(Message Queuing Telemetry Transport,消息队列遥测传输协议)是一种轻量级的发布/订阅协议,专为网络条件苛刻的环境设计,能够在低带宽、不可靠或间歇性通信的场景下稳定运行。目前,MQTT已成为物联网消息通信领域的事实标准协议。

长期稳定更新的攒劲资源: >>>点此立即查看<<<
MQTT的突出特点之一是其提供的三种不同质量的消息服务。在工业物联网场景中,选择MQTT进行数据传输主要基于以下核心优势:
Mosquitto是MQTT协议的一个完整实现,也是业内广泛使用的选择。
近期工作中恰好使用SpringBoot整合Mosquitto,现将完整的搭建过程整理如下。项目基于Maven构建,以下是具体操作步骤。
首先,在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
接下来编写配置文件,将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
创建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);
};
}
}
创建一个消息发送网关,通过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);
}
订阅消息后,需要一个回调处理器来接收并处理消息。这里定义了一个简单的处理器,根据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);
}
}
在实际生产中,通常是在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);
}
在发布消息后,如果设置不当,新的订阅者可能会收到之前发布的历史消息。如果业务上只需要接收即时消息,而不需要历史数据,则需要关注一个参数——retained。
Retained 消息是指在PUBLISH数据包中Retain标识设为1的消息。Broker收到此类消息后,会将其保存下来,当有新的订阅者订阅相应主题时,Broker会立即将该消息推送给订阅者。
因此,在生产者端,只需将retained设置为false,即可避免保留历史消息。至于QoS等级,可根据实际业务需求灵活设定。至此,配置基本完成。
以上就是完整的SpringBoot整合MQTT的实战配置。从协议简介到具体代码实现,覆盖了生产者和消费者的核心环节,希望能为大家提供一个实用的参考。
侠游戏发布此文仅为了传递信息,不代表侠游戏网站认同其观点或证实其描述