最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

SpringBoot實(shí)現(xiàn)MQTT消息發(fā)送和接收方式

 更新時(shí)間:2023年03月11日 10:21:17   作者:LoveDR_1995  
這篇文章主要介紹了SpringBoot實(shí)現(xiàn)MQTT消息發(fā)送和接收方式,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教

Spring integration交互邏輯

對(duì)于發(fā)布者:

1.消息通過消息網(wǎng)關(guān)發(fā)送出去,由 MessageChannel 的實(shí)例 DirectChannel 處理發(fā)送的細(xì)節(jié)。

2.DirectChannel 收到消息后,內(nèi)部通過 MessageHandler 的實(shí)例 MqttPahoMessageHandler 發(fā)送到指定的 Topic。

對(duì)于訂閱者:

1.通過注入 MessageProducerSupport 的實(shí)例 MqttPahoMessageDrivenChannelAdapter,實(shí)現(xiàn)訂閱 Topic 和綁定消息消費(fèi)的 MessageChannel

2.同樣由 MessageChannel 的實(shí)例 DirectChannel 處理消費(fèi)細(xì)節(jié)。

Channel 消息后會(huì)發(fā)送給我們自定義的 MqttInboundMessageHandler 實(shí)例進(jìn)行消費(fèi)。

可以看到整個(gè)處理的流程和前面將的基本一致。Spring Integration 就是抽象出了這么一套消息通信的機(jī)制,具體的通信細(xì)節(jié)由它集成的中間件來決定

1、maven依賴

<!-- https://mvnrepository.com/artifact/org.springframework.boot/spring-boot-starter-integration -->
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-integration</artifactId>
    <version>2.5.1</version>
</dependency>
 
<!-- https://mvnrepository.com/artifact/org.springframework.integration/spring-integration-stream -->
<dependency>
    <groupId>org.springframework.integration</groupId>
    <artifactId>spring-integration-stream</artifactId>
    <version>5.5.5</version>
</dependency>
<!-- https://mvnrepository.com/artifact/org.springframework.integration/spring-integration-mqtt -->
<dependency>
    <groupId>org.springframework.integration</groupId>
    <artifactId>spring-integration-mqtt</artifactId>
    <version>5.5.5</version>
</dependency>

2、yaml配置文件

#mqtt配置
mqtt:
  username: 123
  password: 123
  #MQTT-服務(wù)器連接地址,如果有多個(gè),用逗號(hào)隔開
  url: tcp://127.0.0.1:1883
  #MQTT-連接服務(wù)器默認(rèn)客戶端ID
  client:
    id: ${random.value}
  default:
    #MQTT-默認(rèn)的消息推送主題,實(shí)際可在調(diào)用接口時(shí)指定
    topic: topic,mqtt/test/#
    #連接超時(shí)
  completionTimeout: 3000

3、mqtt生產(chǎn)者消費(fèi)者配置類

import lombok.extern.slf4j.Slf4j;
import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.integration.annotation.IntegrationComponentScan;
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.Message;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.MessageHandler;
import org.springframework.messaging.MessagingException;
 
import java.util.Arrays;
import java.util.List;
 
/**
 * mqtt 推送and接收 消息類
 **/
@Configuration
@IntegrationComponentScan
@Slf4j
public class MqttSenderAndReceiveConfig {
 
    private static final byte[] WILL_DATA;
 
    static {
        WILL_DATA = "offline".getBytes();
    }
 
    @Autowired
    private MqttReceiveHandle mqttReceiveHandle;
 
    @Value("${mqtt.username}")
    private String username;
 
    @Value("${mqtt.password}")
    private String password;
 
    @Value("${mqtt.url}")
    private String hostUrl;
 
    @Value("${mqtt.client.id}")
    private String clientId;
 
    @Value("${mqtt.default.topic}")
    private String defaultTopic;
 
    @Value("${mqtt.completionTimeout}")
    private int completionTimeout;   //連接超時(shí)
 
    /**
     * MQTT連接器選項(xiàng)
     **/
    @Bean(value = "getMqttConnectOptions")
    public MqttConnectOptions getMqttConnectOptions1() {
        MqttConnectOptions mqttConnectOptions = new MqttConnectOptions();
        // 設(shè)置是否清空session,這里如果設(shè)置為false表示服務(wù)器會(huì)保留客戶端的連接記錄,這里設(shè)置為true表示每次連接到服務(wù)器都以新的身份連接
        mqttConnectOptions.setCleanSession(true);
        // 設(shè)置超時(shí)時(shí)間 單位為秒
        mqttConnectOptions.setConnectionTimeout(10);
        mqttConnectOptions.setAutomaticReconnect(true);
        mqttConnectOptions.setUserName(username);
        mqttConnectOptions.setPassword(password.toCharArray());
        mqttConnectOptions.setServerURIs(new String[]{hostUrl});
        // 設(shè)置會(huì)話心跳時(shí)間 單位為秒 服務(wù)器會(huì)每隔1.5*20秒的時(shí)間向客戶端發(fā)送心跳判斷客戶端是否在線,但這個(gè)方法并沒有重連的機(jī)制
        mqttConnectOptions.setKeepAliveInterval(10);
        // 設(shè)置“遺囑”消息的話題,若客戶端與服務(wù)器之間的連接意外中斷,服務(wù)器將發(fā)布客戶端的“遺囑”消息。
        //mqttConnectOptions.setWill("willTopic", WILL_DATA, 2, false);
        return mqttConnectOptions;
    }
 
    /**
     * MQTT工廠
     **/
    @Bean
    public MqttPahoClientFactory mqttClientFactory() {
        DefaultMqttPahoClientFactory factory = new DefaultMqttPahoClientFactory();
        factory.setConnectionOptions(getMqttConnectOptions1());
        return factory;
    }
 
    /**
     * MQTT信息通道(生產(chǎn)者)
     **/
    @Bean
    public MessageChannel mqttOutboundChannel() {
        return new DirectChannel();
    }
 
    /**
     * MQTT消息處理器(生產(chǎn)者)
     **/
    @Bean
    @ServiceActivator(inputChannel = "mqttOutboundChannel")
    public MessageHandler mqttOutbound() {
        MqttPahoMessageHandler messageHandler = new MqttPahoMessageHandler(clientId + "_producer", mqttClientFactory());
        messageHandler.setAsync(true);
        messageHandler.setDefaultTopic(defaultTopic);
        messageHandler.setAsyncEvents(true); // 消息發(fā)送和傳輸完成會(huì)有異步的通知回調(diào)
        //設(shè)置轉(zhuǎn)換器 發(fā)送bytes數(shù)據(jù)
        DefaultPahoMessageConverter converter = new DefaultPahoMessageConverter();
        converter.setPayloadAsBytes(true);
        return messageHandler;
    }
 
    /**
     * 配置client,監(jiān)聽的topic
     * MQTT消息訂閱綁定(消費(fèi)者)
     **/
    @Bean
    public MessageProducer inbound() {
        List<String> topicList = Arrays.asList(defaultTopic.trim().split(","));
        String[] topics = new String[topicList.size()];
        topicList.toArray(topics);
        MqttPahoMessageDrivenChannelAdapter adapter =
                new MqttPahoMessageDrivenChannelAdapter(clientId + "_consumer", mqttClientFactory(), topics);
        adapter.setCompletionTimeout(completionTimeout);
        DefaultPahoMessageConverter converter = new DefaultPahoMessageConverter();
        converter.setPayloadAsBytes(true);
        adapter.setConverter(converter);
        adapter.setQos(2);
        adapter.setOutputChannel(mqttInputChannel());
        return adapter;
    }
 
    /**
     * MQTT信息通道(消費(fèi)者)
     **/
    @Bean
    public MessageChannel mqttInputChannel() {
        return new DirectChannel();
    }
 
    /**
     * MQTT消息處理器(消費(fèi)者)
     **/
    @Bean
    @ServiceActivator(inputChannel = "mqttInputChannel")
    public MessageHandler handler() {
        return new MessageHandler() {
            @Override
            public void handleMessage(Message<?> message) throws MessagingException {
                //處理接收消息
                mqttReceiveHandle.handle(message);
            }
        };
    }
}

4、消息處理類 

/**
 * mqtt客戶端消息處理類
 **/
@Slf4j
@Component
public class MqttReceiveHandle {
 
    public void handle(Message<?> message) {
        log.info("收到訂閱消息: {}", message);
        String topic = message.getHeaders().get(MqttHeaders.RECEIVED_TOPIC).toString();
        log.info("消息主題:{}", topic);
        Object payLoad = message.getPayload();
        byte[] data = (byte[]) payLoad;
        Packet packet = Packet.parse(data);
        log.info("發(fā)送的Packet數(shù)據(jù){}", JSON.toJSONString(packet));
 
    }
}

5、mqtt發(fā)送接口 

import org.springframework.integration.annotation.MessagingGateway;
import org.springframework.integration.mqtt.support.MqttHeaders;
import org.springframework.messaging.handler.annotation.Header;
 
/**
 * mqtt發(fā)送消息
 * (defaultRequestChannel = "mqttOutboundChannel" 對(duì)應(yīng)config配置)
 * **/
@MessagingGateway(defaultRequestChannel = "mqttOutboundChannel")
public interface MqttGateway {
 
    /**
     * 發(fā)送信息到MQTT服務(wù)器
     *
     * @param
     */
    void sendToMqttObject(@Header(MqttHeaders.TOPIC) String topic, byte[] payload);
 
    /**
     * 發(fā)送信息到MQTT服務(wù)器
     *
     * @param topic 主題
     * @param payload 消息主體
     */
    void sendToMqtt(@Header(MqttHeaders.TOPIC) String topic, String payload);
 
    /**
     * 發(fā)送信息到MQTT服務(wù)器
     *
     * @param topic 主題
     * @param qos 對(duì)消息處理的幾種機(jī)制。
     * 0 表示的是訂閱者沒收到消息不會(huì)再次發(fā)送,消息會(huì)丟失。
     * 1 表示的是會(huì)嘗試重試,一直到接收到消息,但這種情況可能導(dǎo)致訂閱者收到多次重復(fù)消息。
     * 2 多了一次去重的動(dòng)作,確保訂閱者收到的消息有一次。
     * @param payload 消息主體
     */
    void sendToMqtt(@Header(MqttHeaders.TOPIC) String topic, @Header(MqttHeaders.QOS) int qos, String payload);
 
    /**
     * 發(fā)送信息到MQTT服務(wù)器
     *
     * @param topic 主題
     * @param payload 消息主體
     */
    void sendToMqtt(@Header(MqttHeaders.TOPIC) String topic, Object payload);
 
    /**
     * 發(fā)送信息到MQTT服務(wù)器
     *
     * @param topic 主題
     * @param payload 消息主體
     */
    void sendToMqtt(@Header(MqttHeaders.TOPIC) String topic, byte[] payload);
}

6、mqtt事件監(jiān)聽類 

import lombok.extern.slf4j.Slf4j;
import org.springframework.context.event.EventListener;
import org.springframework.integration.mqtt.event.MqttConnectionFailedEvent;
import org.springframework.integration.mqtt.event.MqttMessageDeliveredEvent;
import org.springframework.integration.mqtt.event.MqttMessageSentEvent;
import org.springframework.integration.mqtt.event.MqttSubscribedEvent;
import org.springframework.stereotype.Component;
 
@Slf4j
@Component
public class MqttListener {
    /**
     * 連接失敗的事件通知
     * @param mqttConnectionFailedEvent
     */
    @EventListener(classes = MqttConnectionFailedEvent.class)
    public void listenerAction(MqttConnectionFailedEvent mqttConnectionFailedEvent) {
        log.info("連接失敗的事件通知");
    }
 
    /**
     * 已發(fā)送的事件通知
     * @param mqttMessageSentEvent
     */
    @EventListener(classes = MqttMessageSentEvent.class)
    public void listenerAction(MqttMessageSentEvent mqttMessageSentEvent) {
        log.info("已發(fā)送的事件通知");
    }
 
    /**
     * 已傳輸完成的事件通知
     * 1.QOS == 0,發(fā)送消息后會(huì)即可進(jìn)行此事件回調(diào),因?yàn)椴恍枰却貓?zhí)
     * 2.QOS == 1,發(fā)送消息后會(huì)等待ACK回執(zhí),ACK回執(zhí)后會(huì)進(jìn)行此事件通知
     * 3.QOS == 2,發(fā)送消息后會(huì)等待PubRECV回執(zhí),知道收到PubCOMP后會(huì)進(jìn)行此事件通知
     * @param mqttMessageDeliveredEvent
     */
    @EventListener(classes = MqttMessageDeliveredEvent.class)
    public void listenerAction(MqttMessageDeliveredEvent mqttMessageDeliveredEvent) {
        log.info("已傳輸完成的事件通知");
    }
 
    /**
     * 消息訂閱的事件通知
     * @param mqttSubscribedEvent
     */
    @EventListener(classes = MqttSubscribedEvent.class)
    public void listenerAction(MqttSubscribedEvent mqttSubscribedEvent) {
        log.info("消息訂閱的事件通知");
    }
}

7、接口測(cè)試

@Resource
    private MqttGateway mqttGateway;
    /**
     * sendData 消息
     * topic 訂閱主題
     **/
    @RequestMapping(value = "/sendMqtt",method = RequestMethod.POST)
    public String sendMqtt(String sendData, String topic) {
        MqttMessage mqttMessage = new MqttMessage();
        mqttGateway.sendToMqtt(topic, sendData);
        //mqttGateway.sendToMqttObject(topic, sendData.getBytes());
        return "OK";
    }

總結(jié)

以上為個(gè)人經(jīng)驗(yàn),希望能給大家一個(gè)參考,也希望大家多多支持腳本之家。

相關(guān)文章

  • spring boot配置ssl(多cer格式)超詳細(xì)教程

    spring boot配置ssl(多cer格式)超詳細(xì)教程

    這篇文章主要介紹了spring boot配置ssl(多cer格式)超詳細(xì)教程,本文通過圖文并茂的形式給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友參考下吧
    2023-11-11
  • 詳解Java8?CompletableFuture的并行處理用法

    詳解Java8?CompletableFuture的并行處理用法

    Java8中有一個(gè)工具非常有用,那就是CompletableFuture,本章主要講解CompletableFuture的并行處理用法,感興趣的小伙伴可以了解一下
    2022-04-04
  • swagger注解@ApiModelProperty失效情況的解決

    swagger注解@ApiModelProperty失效情況的解決

    這篇文章主要介紹了swagger注解@ApiModelProperty失效情況的解決方案,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2022-06-06
  • springAOP完整實(shí)現(xiàn)過程

    springAOP完整實(shí)現(xiàn)過程

    當(dāng)你調(diào)用SimpleService類的doSomething方法時(shí),上述的PerformanceAspect會(huì)自動(dòng)攔截此調(diào)用,并且記錄該方法的執(zhí)行時(shí)間,這樣你就完成了一個(gè)針對(duì)Spring的AOP入門級(jí)案例,感興趣的朋友一起看看吧
    2024-02-02
  • 使用MyBatis-Generator如何自動(dòng)生成映射文件

    使用MyBatis-Generator如何自動(dòng)生成映射文件

    這篇文章主要介紹了使用MyBatis-Generator如何自動(dòng)生成映射文件,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2022-02-02
  • MyBatis Plus關(guān)閉SQL日志打印的方法

    MyBatis Plus關(guān)閉SQL日志打印的方法

    這篇文章主要介紹了MyBatis-Plus如何關(guān)閉SQL日志打印,文中通過圖文結(jié)合講解的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2024-02-02
  • Vert.x學(xué)習(xí)之Resilience4j原理與用法解讀

    Vert.x學(xué)習(xí)之Resilience4j原理與用法解讀

    Resilience4j通過獨(dú)立線程和異步控制機(jī)制,為Vert.x提供了強(qiáng)大的容錯(cuò)能力,包括斷路器、重試、限流和超時(shí)控制,確保系統(tǒng)在面對(duì)故障時(shí)仍能穩(wěn)定運(yùn)行
    2026-01-01
  • SpringBoot調(diào)用WebService接口的實(shí)現(xiàn)示例

    SpringBoot調(diào)用WebService接口的實(shí)現(xiàn)示例

    本文主要介紹了SpringBoot調(diào)用WebService接口的實(shí)現(xiàn)示例,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2025-03-03
  • 基于 SpringBoot 實(shí)現(xiàn) MySQL 讀寫分離的問題

    基于 SpringBoot 實(shí)現(xiàn) MySQL 讀寫分離的問題

    這篇文章主要介紹了基于 SpringBoot 實(shí)現(xiàn) MySQL 讀寫分離的問題,本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2021-02-02
  • Spring?Boot?打包如何將依賴全部打進(jìn)去

    Spring?Boot?打包如何將依賴全部打進(jìn)去

    這篇文章主要介紹了Spring?Boot?打包如何將依賴全部打進(jìn)去,在pom.xml中引入插件,需要在項(xiàng)目的pom.xml文件中,添加?Maven?插件??spring-boot-maven-plugin,本文結(jié)合實(shí)例代碼介紹的非常詳細(xì),需要的朋友可以參考下
    2023-09-09

最新評(píng)論

铜鼓县| 玉田县| 紫金县| 辉南县| 丰顺县| 孙吴县| 鸡东县| 阳曲县| 望都县| 双鸭山市| 宜都市| 兴安盟| 金华市| 封丘县| 浙江省| 株洲市| 定州市| 旬邑县| 广南县| 九龙县| 平潭县| 高清| 玛纳斯县| 苍溪县| 新蔡县| 交城县| 浠水县| 云林县| 仙桃市| 奇台县| 丹东市| 镇平县| 惠东县| 礼泉县| 手机| 安溪县| 怀安县| 泰来县| 桃江县| 南溪县| 昌平区|