java代碼mqtt接收發(fā)送消息方式
java代碼mqtt接收發(fā)送消息
mqtt消息第一用到不是太熟悉所以寫(xiě)一篇文章鞏固一下。
前提是你已經(jīng)把mqtt已經(jīng)安裝好,并且啟動(dòng)好了。
首先我們需要兩部分代碼。
所需依賴(lài)
<!-- mqtt -->
<dependency>
<groupId>org.springframework.integration</groupId>
<artifactId>spring-integration-stream</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.integration</groupId>
<artifactId>spring-integration-mqtt</artifactId>
</dependency>連接mqtt部分的代碼塊,因?yàn)槲也恍枰l(fā)送消息所以把發(fā)送消息給注釋掉了。
package mqttclient.util;
import lombok.extern.slf4j.Slf4j;
import mqttclient.callback.MqttMessageCallback2;
import org.eclipse.paho.client.mqttv3.MqttClient;
import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
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 MqttClientUtil2 {
private String username;
private String password;
@Value("tcp://127.0.0.1:1883")//這個(gè)是安裝mqtt的ip以及端口,1883是mqtt默認(rèn)端口
private String host;
@Value("CYT")//這個(gè)隨便寫(xiě)但是是唯一的。
private String clientId;
@Value("cyt/#")這個(gè)是mqtt發(fā)送消息的咱們要訂閱的topic,cyt/#代表以cyt/開(kāi)始的所有topic都接收
private String topic;
@Value("${mqtt.connection.timeout}")//IOT_MQTT_Yield會(huì)block住timeout的時(shí)間去嘗試接收數(shù)據(jù),直到timeout才會(huì)退出??梢詫?xiě)在這里也可以寫(xiě)在yml配置文件中
private int timeOut;
@Value("${mqtt.keep.alive.interval}")
private int interval;
@Autowired
private MqttMessageCallback2 mqttMessageCallback2;
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();
//設(shè)置回調(diào)函數(shù)
mqttClient.setCallback(mqttMessageCallback2);
//鏈接mqtt
mqttClient.connect(mqttConnectOptions);
//訂閱消息
mqttClient.subscribe(topic,2);
}catch (Exception e){
log.error("mqtt服務(wù)鏈接異常!");
e.printStackTrace();
}
}
/**
* 設(shè)置鏈接對(duì)象信息
* setCleanSession true 斷開(kāi)鏈接即清楚會(huì)話 false 保留鏈接信息 離線還會(huì)繼續(xù)發(fā)消息
* @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鏈接狀態(tài)
* @return
*/
private boolean isConnect(){
if(Objects.isNull(this.mqttClient)){
return false;
}
return mqttClient.isConnected();
}
/**
* 設(shè)置重連
* @throws Exception
*/
private void reConnect() throws Exception{
if(Objects.nonNull(this.mqttClient)){
log.info("mqtt 服務(wù)已重新鏈接...");
this.mqttClient.connect(this.mqttConnectOptions);
}
}
/**
* 斷開(kāi)鏈接
* @throws Exception
*/
private void closeConnect() throws Exception{
if(Objects.nonNull(this.mqttClient)){
log.info("mqtt 服務(wù)已斷開(kāi)鏈接...");
this.mqttClient.disconnect();
}
}
/* *//**
* 發(fā)布消息
* @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("消息發(fā)送成功---->{}",message);
}
}catch(Exception e){
log.error("消息發(fā)送異常",e);
}
}
}else{
reConnect();
}
}*/
}接收消息部分
package mqttclient.callback;
import lombok.extern.slf4j.Slf4j;
import mqttclient.util.ParsingData2;
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;
import java.util.List;
@Component
@Slf4j
public class MqttMessageCallback2 implements MqttCallback {
/**
* 鏈接丟失時(shí)處理
* @param throwable
*/
@Override
public void connectionLost(Throwable throwable) {
//可以做重連 或者 其他業(yè)務(wù)處理
}
@Override
public void messageArrived(String topic, MqttMessage mqttMessage) throws Exception {
System.out.println("接收到消息topic---->{}"+topic);
System.out.println("接收到消息topic---->{}"+mqttMessage);
log.info("接收到消息質(zhì)量qos---->{}",mqttMessage.getQos());
System.out.println("接收到消息質(zhì)量qos---->{}"+mqttMessage.getQos());
log.info("接收到消息具體信息---->{}",new String(mqttMessage.getPayload()));
System.out.println("接收到消息具體信息---->{}"+mqttMessage.getPayload());
//結(jié)合業(yè)務(wù) 編寫(xiě)具體信息即可
}
@Override
public void deliveryComplete(IMqttDeliveryToken iMqttDeliveryToken) {
}
}這個(gè)兩個(gè)寫(xiě)完之后只要有數(shù)據(jù)發(fā)送過(guò)來(lái),這邊會(huì)自動(dòng)進(jìn)行接收打印。
是用mqtt網(wǎng)頁(yè)版圖形化界面進(jìn)行模擬數(shù)據(jù)發(fā)送。
安裝mqtt后打開(kāi)此網(wǎng)站:http://localhost:18083/
默認(rèn)賬號(hào)是:admin / public
登錄后這邊可以設(shè)置中文:

模擬發(fā)送:這幾個(gè)地方不用改動(dòng)但是一定要點(diǎn)擊綠色的連接才可以,進(jìn)行發(fā)送。

需要修改的部分是:

然后點(diǎn)擊發(fā)送就可以收到信息了。
總結(jié)
以上為個(gè)人經(jīng)驗(yàn),希望能給大家一個(gè)參考,也希望大家多多支持腳本之家。
相關(guān)文章
Spring?Security重寫(xiě)AuthenticationManager實(shí)現(xiàn)賬號(hào)密碼登錄或者手機(jī)號(hào)碼登錄
本文主要介紹了Spring?Security重寫(xiě)AuthenticationManager實(shí)現(xiàn)賬號(hào)密碼登錄或者手機(jī)號(hào)碼登錄,文中通過(guò)示例代碼介紹的非常詳細(xì),需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧2025-08-08
SpringBoot啟動(dòng)后自動(dòng)執(zhí)行方法的各種方式對(duì)比
這篇文章主要為大家詳細(xì)介紹了SpringBoot啟動(dòng)后自動(dòng)執(zhí)行方法的各種方式和性能對(duì)比,文中的示例代碼講解詳細(xì),感興趣的小伙伴可以參考一下2025-04-04
Java新特性中Preview功能如何運(yùn)行調(diào)試詳解
這篇文章主要為大家介紹了Java新特性中Preview功能如何運(yùn)行調(diào)試詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪2023-10-10

