SpringBoot整合MQTT協(xié)議實現(xiàn)消息訂閱與發(fā)布功能
1、相關(guān)依賴 pom.xml文件
<dependency>
<groupId>org.springframework.integration</groupId>
<artifactId>spring-integration-mqtt</artifactId>
</dependency>2、配置文件 application.yml
這里的訂閱主題可不要,我這里用于啟動的時候就訂閱固定主題。適用主題固定的場景。
# MQTT服務(wù)地址,端口默認(rèn)1883 mqtt-broker-url: tcp://127.0.0.1:1883 # 用戶名 mqtt-username: admin # 密碼 mqtt-password: public # 訂閱主題(可以多個) mqtt-default-topic: mqtt/topic_test # 客戶端Id mqtt-clientId: can
3、MQTT配置類
用于配置項目啟動時就連接MQTT。
@Component
public class MqttConfig {
@Resource
private MqttPushClient mqttPushClient;
@Resource
private MqttSubClient mqttSubClient;
/**
* 用戶名
*/
@Value("${mqtt-username}")
private String username;
/**
* 密碼
*/
@Value("${mqtt-password}")
private String password;
/**
* 連接地址
*/
@Value("${mqtt-broker-url}")
private String hostUrl;
/**
* 客戶Id
*/
@Value("${mqtt-clientId}")
private String clientId;
/**
* 默認(rèn)連接話題,多個的話用逗號隔開
*/
@Value("${mqtt-default-topic}")
private String defaultTopic;
/**
* 超時時間
*/
private int timeout = 100;
/**
* 保持連接數(shù)
*/
private int keepalive = 60;
/**
* 連接至mqtt服務(wù)器,獲取mqtt連接
*
* @return MqttPushClient
*/
@Bean
public MqttPushClient getMqttPushClient() {
// 連接至mqtt服務(wù)器,獲取mqtt連接
mqttPushClient.connect(hostUrl, clientId, username, password, timeout, keepalive);
// 訂閱默認(rèn)主題
mqttSubClient.subScribeDataPublishTopic(defaultTopic);
return mqttPushClient;
}
}4、發(fā)布連接類
連接MQTT的方法、發(fā)布消息的方法。
@Slf4j
@Component
public class MqttPushClient {
private static final Logger logger = LoggerFactory.getLogger(MqttPushClient.class);
@Autowired
private PushCallback pushCallback;
@Getter
private static MqttClient client;
public static void setClient(MqttClient client) {
MqttPushClient.client = client;
}
/**
* 連接
* @param host mqtt://127.0.0.1:1883
* @param clientId can
* @param username admin
* @param password password
* @param timeout 100
* @param keepalive 60
*/
public void connect(String host, String clientId, String username, String password, int timeout, int keepalive) {
MqttClient client;
try {
client = new MqttClient(host, clientId, new MemoryPersistence());
MqttConnectOptions options = new MqttConnectOptions();
options.setCleanSession(true);
options.setUserName(username);
options.setPassword(password.toCharArray());
options.setConnectionTimeout(timeout);
options.setKeepAliveInterval(keepalive);
// automaticReconnect 為 true 表示斷線自動重連,但僅僅只是重新連接,并不訂閱主題;在 connectComplete 回調(diào)函數(shù)重新訂閱
options.setAutomaticReconnect(true);
MqttPushClient.setClient(client);
try {
//設(shè)置回調(diào)類
client.setCallback(pushCallback);
IMqttToken iMqttToken = client.connectWithResult(options);
boolean complete = iMqttToken.isComplete();
log.error("MQTT連接{}", complete ? "成功" : "失敗");
} catch (Exception e) {
logger.error(e.getMessage());
e.printStackTrace();
}
} catch (Exception e) {
logger.error(e.getMessage());
e.printStackTrace();
}
}
/**
* 關(guān)閉MQTT連接
*/
public void close() throws MqttException {
client.disconnect();
client.close();
}
/**
* 發(fā)布,默認(rèn)qos為0,非持久化
*
* @param topic 主題名
* @param pushMessage 消息
*/
public void publish(String topic, String pushMessage) {
publish(0, false, topic, pushMessage);
}
/**
* 發(fā)布
* QoS 0:消息最多傳送一次。如果當(dāng)前客戶端不可用,它將丟失這條消息。
* QoS 1:消息至少傳送一次。
* QoS 2:消息只傳送一次。
* @param qos
* @param retained
* @param topic
* @param pushMessage
*/
public void publish(int qos, boolean retained, String topic, String pushMessage) {
MqttMessage message = new MqttMessage();
message.setQos(qos);
message.setRetained(retained);
message.setPayload(pushMessage.getBytes());
MqttTopic mTopic = MqttPushClient.getClient().getTopic(topic);
// MQTT主題不存在
if (null == mTopic) return;
try {
mTopic.publish(message);
} catch (Exception e) {
log.error("MQTT發(fā)送消息異常:", e);
e.printStackTrace();
}
}
}5、訂閱類
用于訂閱某個或多個主題、取消訂閱某個或者多個主題。
@Slf4j
@Component
public class MqttSubClient {
private static final Logger logger = LoggerFactory.getLogger(MqttSubClient.class);
// 訂閱多個主題以逗號分開
public void subScribeDataPublishTopic(String defaultTopic) {
//訂閱test_queue主題
String[] mqttTopic = defaultTopic.split(",");
for (String s : mqttTopic) {
//訂閱主題
subscribe(s, 0);
}
}
/**
* 訂閱某個主題,qos默認(rèn)為0
*
* @param topic 主題
*/
public void subscribe(String topic) {
subscribe(topic, 0);
}
/**
* 訂閱某個主題
*
* @param topic 主題名
* @param qos qos
*/
public void subscribe(String topic, int qos) {
try {
MqttClient client = MqttPushClient.getClient();
if (client == null) {
return;
}
client.subscribe(topic, qos);
log.error("MQTT訂閱主題:{}", topic);
} catch (MqttException e) {
logger.error(e.getMessage());
e.printStackTrace();
}
}
/**
* 取消訂閱某個主題
* @param topic 要取消訂閱的主題名
*/
public void unsubscribe(String topic) {
try {
MqttClient client = MqttPushClient.getClient();
if (client == null || !client.isConnected()) {
return;
}
client.unsubscribe(topic); // 取消訂閱
log.error("MQTT取消訂閱主題: {}", topic);
} catch (MqttException e) {
log.error("取消訂閱失敗: {}", e.getMessage());
e.printStackTrace();
}
}
/**
* 批量取消訂閱多個主題
* @param topics 主題數(shù)組
*/
public void unsubscribe(String[] topics) {
try {
MqttClient client = MqttPushClient.getClient();
if (client == null || !client.isConnected()) {
return;
}
client.unsubscribe(topics); // 取消訂閱多個主題
log.error("MQTT取消訂閱主題: {}", Arrays.toString(topics));
} catch (MqttException e) {
log.error("取消訂閱失敗: {}", e.getMessage());
e.printStackTrace();
}
}
}
6、回調(diào)類
處理MQTT連接斷開重連、訂閱主題接收的消息處理。
@Slf4j
@Component
public class PushCallback implements MqttCallback {
@Resource
@Lazy
private MqttPushClient mqttPushClient;
/**
* 連接丟失后,一般在這里面進行重連(重連的邏輯需要自己處理)
* @param cause .
*/
@Override
public void connectionLost(Throwable cause) {
log.error("MQTT連接斷開,正在重連:" + cause);
}
/**
* 發(fā)送消息,消息到達后處理方法
* @param token .
*/
@Override
public void deliveryComplete(IMqttDeliveryToken token) {
log.error("deliveryComplete---------{}", token.isComplete());
}
/**
* 訂閱主題接收到消息處理方法
* @param topic 主題
* @param message 消息
*/
@Override
public void messageArrived(String topic, MqttMessage message) {
// 訂閱主題后得到的消息會執(zhí)行到這里面,這里在控制臺有輸出
log.error("MQTT接收消息主題 : {}", topic);
log.error("MQTT接收消息Qos : {}", message.getQos());
log.error("MQTT接收消息內(nèi)容 : {}", message);
}
}7、啟動后,進入EMQX管理頁面
程序允許打印連接成功,去EMQX管理頁面查看。

EMQX管理頁面這里有所有的主題列表。

包括客戶端訂閱的主題。

8、通過接口給主題發(fā)送消息
@RestController
@Slf4j
@RequestMapping("/api")
public class ApiController {
@Resource
private MqttPushClient mqttPushClient;
@GetMapping("/test")
public String getVersions(@RequestParam String topic, @RequestParam String message) {
mqttPushClient.publish(topic, message);
return "ok";
}
}瀏覽器直接調(diào)用,topic:配置文件里面訂閱的主題。message:你想發(fā)送給主題的消息。

控制臺日志打?。喊l(fā)送成功,并且接收到了主題發(fā)送的消息。

到此這篇關(guān)于SpringBoot整合MQTT協(xié)議實現(xiàn)消息訂閱與發(fā)布功能的文章就介紹到這了,更多相關(guān)SpringBoot整合MQTT訂閱與發(fā)布內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
詳解如何在spring boot中使用spring security防止CSRF攻擊
這篇文章主要介紹了詳解如何在spring boot中使用spring security防止CSRF攻擊,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧2018-05-05
Java內(nèi)部類_動力節(jié)點Java學(xué)院整理
內(nèi)部類是指在一個外部類的內(nèi)部再定義一個類。下面通過本文給大家java內(nèi)部類的使用小結(jié),需要的朋友參考下吧2017-04-04
Java使用反射和動態(tài)代理實現(xiàn)一個View注解綁定庫
這篇文章主要介紹了Java使用反射和動態(tài)代理實現(xiàn)一個View注解綁定庫,代碼簡潔,使用簡單,擴展性強,結(jié)合實例代碼給大家介紹的非常詳細,需要的朋友可以參考下2022-05-05
java獲取兩個數(shù)組中不同數(shù)據(jù)的方法
這篇文章主要介紹了java獲取兩個數(shù)組中不同數(shù)據(jù)的方法,實例分析了java操作數(shù)組的技巧,非常具有實用價值,需要的朋友可以參考下2015-03-03
Java接口回調(diào)和方法回調(diào)的簡單實現(xiàn)步驟
這篇文章主要介紹了Java接口回調(diào)和方法回調(diào)的相關(guān)資料,接口回調(diào)是一種設(shè)計模式,實現(xiàn)三方解耦,調(diào)用者提供接口實現(xiàn),文中通過代碼介紹的非常詳細,需要的朋友可以參考下2025-03-03
Spring?Boot?JWT登錄授權(quán)使用(無感刷新)
JWT作為一種輕量級的身份認(rèn)證與授權(quán)方案,憑借其無狀態(tài)、可跨域、易于擴展的特性,成為?Spring?Boot?項目中實現(xiàn)認(rèn)證授權(quán)的主流選擇,本文詳解Spring?Boot整合?JWT?實現(xiàn)登錄認(rèn)證與授權(quán)的全流程,助力開發(fā)者快速搭建安全可靠的認(rèn)證授權(quán)體系,感興趣的可以了解一下2026-04-04
SpringMVC之RequestContextHolder詳細解析
這篇文章主要介紹了SpringMVC之RequestContextHolder詳細解析,正常來說在service層是沒有request的,然而直接從controlller傳過來的話解決方法太粗暴,后來發(fā)現(xiàn)了SpringMVC提供的RequestContextHolder,需要的朋友可以參考下2023-11-11
Java 創(chuàng)建兩個線程模擬對話并交替輸出實現(xiàn)解析
這篇文章主要介紹了Java 創(chuàng)建兩個線程模擬對話并交替輸出實現(xiàn)解析,文中通過示例代碼介紹的非常詳細,對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友可以參考下2019-10-10

