springboot整合mqtt
·
1、引入依赖
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-integration</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>
2、配置项
mqtt:
serverURIs: tcp://127.0.0.1:1883 #mqtt服务地址
username: abc
password: abc
keepAliveInterval: 3 #心跳 单位 秒
connectionTimeout: 5 #连接超时时间 单位 秒
cleanSession: false #是否清除session
async: true #是否异步发送
producer: #生产者配置
defaultQos: 1 #消息质量 0:最多一次传输(可能丢包) 1:至少一次传输,(可能重包) 2:保证只有一次传输
defaultRetained: false #是否保留消息
clientId: server_outbound #生产者连接id
defaultTopic: notification #默认订阅
consumer: #消费者配置
defaultTopic: notification #默认订阅
clientId: server_inbound #消费者连接id
3、MQTT配置类
@Configuration
@Slf4j
public class MqttConfig {
@Value("${mqtt.serverURIs}")
private String serverURIs;
@Value("${mqtt.username}")
private String username;
@Value("${mqtt.password}")
private String password;
@Value("${mqtt.producer.clientId}")
private String producerClientId;
@Value("${mqtt.consumer.clientId}")
private String consumerClientId;
//心跳频率
@Value("${mqtt.keepAliveInterval}")
private int keepAliveInterval;
//超时时间
@Value("${mqtt.connectionTimeout}")
private int connectionTimeout;
//生产者默认订阅
@Value("${mqtt.producer.defaultTopic}")
private String producerDefaultTopic;
//消费者默认订阅
@Value("${mqtt.consumer.defaultTopic}")
private String consumerDefaultTopic;
//默认消息质量
@Value("${mqtt.producer.defaultQos}")
private int defaultProducerQos;
//默认订阅
@Value("${mqtt.producer.defaultRetained}")
private boolean defaultRetained;
//接收消息管道
public static final String INBOUND_CHANNEL = "server_inbound";
/**
* Mqtt配置项
* @return {@link org.eclipse.paho.client.mqttv3.MqttConnectOptions}
*/
@Bean
public MqttConnectOptions getMqttConnectOptions() {
MqttConnectOptions options = new MqttConnectOptions();
// 设置是否清空session,这里如果设置为false表示服务器会保留客户端的连接记录,
// 这里设置为true表示每次连接到服务器都以新的身份连接
options.setCleanSession(true);
// 用户名
options.setUserName(username);
// 密码
options.setPassword(password.toCharArray());
// mqtt服务url
options.setServerURIs(StringUtils.split(serverURIs, ","));
// 设置超时时间 单位为秒
options.setConnectionTimeout(connectionTimeout);
// 设置会话心跳时间 单位为秒 服务器会每隔1.5*keepAliveInterval秒的时间向客户端发送心跳判断客户端是否在线,但这个方法并没有重连的机制
options.setKeepAliveInterval(keepAliveInterval);
// 设置“遗嘱”消息的话题,若客户端与服务器之间的连接意外中断,服务器将发布客户端的“遗嘱”消息。
//options.setWill("willTopic", WILL_DATA, 2, false);
return options;
}
/**
* MQTT客户端
* @return {@link org.springframework.integration.mqtt.core.MqttPahoClientFactory}
*/
@Bean
public MqttPahoClientFactory mqttClientFactory() {
DefaultMqttPahoClientFactory factory = new DefaultMqttPahoClientFactory();
factory.setConnectionOptions(getMqttConnectOptions());
return factory;
}
/**
* mqtt消息管道(生产者)
* @return {@link org.springframework.messaging.MessageChannel}
*/
@Bean(name = IMqttSender.OUTBOUND_CHANNEL)
public MessageChannel outboundChannel() {
return new DirectChannel();
}
/**
* mqtt消息处理器(生产者)
* @return {@link org.springframework.messaging.MessageHandler}
*/
@Bean
@ServiceActivator(inputChannel = IMqttSender.OUTBOUND_CHANNEL)
public MessageHandler getMqttProducer() {
MqttPahoMessageHandler messageHandler = new MqttPahoMessageHandler(producerClientId, mqttClientFactory());
messageHandler.setAsync(true);
messageHandler.setDefaultTopic(producerDefaultTopic);
//messageHandler.setDefaultRetained(defaultRetained);
messageHandler.setDefaultQos(defaultProducerQos);
return messageHandler;
}
/**
* 消费者 如果不需要订阅消息,可不用配置下面
*/
/**
* MQTT信息通道(消费者)
* @return {@link org.springframework.messaging.MessageChannel}
*/
@Bean(name = INBOUND_CHANNEL)
public MessageChannel mqttInboundChannel() {
return new DirectChannel();
}
/**
* MQTT消息订阅绑定(消费者)
* @return {@link org.springframework.integration.core.MessageProducer}
*/
@Bean
public MessageProducer inbound() {
// 可以同时消费(订阅)多个Topic
MqttPahoMessageDrivenChannelAdapter adapter =
new MqttPahoMessageDrivenChannelAdapter(
consumerClientId, mqttClientFactory(),
StringUtils.split(consumerDefaultTopic, ","));
adapter.setCompletionTimeout(5000);
adapter.setConverter(new DefaultPahoMessageConverter());
adapter.setQos(1);
// 设置订阅通道
adapter.setOutputChannel(mqttInboundChannel());
return adapter;
}
/**
* MQTT消息处理器(消费者) 订阅的消息将会在这里打印
*
* @return {@link org.springframework.messaging.MessageHandler}
*/
@Bean
@ServiceActivator(inputChannel = INBOUND_CHANNEL)
public MessageHandler handler() {
return new MessageHandler() {
@Override
public void handleMessage(Message<?> message) throws MessagingException {
log.info("--------------------接收到订阅消息--------------------");
log.info("{}", message.getPayload());
}
};
}
}
4、MQTT消息发送接口
@Component
@MessagingGateway(defaultRequestChannel = IMqttSender.OUTBOUND_CHANNEL)
public interface IMqttSender {
//发送消息管道
String OUTBOUND_CHANNEL = "server_outbound";
/**
* 发送信息到MQTT服务器
* @param data 发送的文本
*/
void sendToMqtt(String data);
/**
* 发送信息到MQTT服务器
* @param topic 主题
* @param payload 消息主体
*/
void sendToMqtt(@Header(MqttHeaders.TOPIC) String topic,
String payload);
/**
* 发送信息到MQTT服务器
* @param topic 主题
* @param qos 对消息处理的几种机制。<br> 0 表示的是订阅者没收到消息不会再次发送,消息会丢失。<br>
* 1 表示的是会尝试重试,一直到接收到消息,但这种情况可能导致订阅者收到多次重复消息。<br>
* 2 多了一次去重的动作,确保订阅者收到的消息有一次。
* @param payload 消息主体
*/
void sendToMqtt(@Header(MqttHeaders.TOPIC) String topic,
@Header(MqttHeaders.QOS) int qos,
String payload);
}
注意事项:@MessagingGateway注解不会被@ComponentScan当做普通的组件扫描。放在Application同级目录及子目录时,会被默认扫描。如果需要将IMqttSender放在Application类的目录外,需要在Application类中配置@IntegrationComponentScan进行扫描。
5、发送消息到MQTT服务
@ActiveProfiles("test")
@SpringBootTest
public class MqttTest {
@Resource
private IMqttSender mqttSender;
@Test
void contextLoads(){
mqttSender.sendToMqtt("测试发布消息");//默认topic 默认qos
mqttSender.sendToMqtt("testTopic","测试发布消息");//自定义topic 默认qos
mqttSender.sendToMqtt("testTopic",1,"测试发布消息");//自定义topic 自定义qos
}
}
更多推荐
所有评论(0)