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

SpringBoot集成MQTT示例詳解

 更新時間:2022年07月22日 14:07:00   作者:bluesbruce  
這篇文章主要為大家介紹了SpringBoot集成MQTT示例詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪

引言

特別提醒: 文中提到的MQTT服務器Apache-Apollo,現(xiàn)在已經(jīng)不維護。但是客戶端的寫法是通用的。目前我常用的是RabbitMQ加mqtt插件。

MQTT

MQTT(消息隊列遙測傳輸)是ISO標準(ISO/IEC PRF 20922)下基于發(fā)布/訂閱范式的消息協(xié)議。它工作在 TCP/IP協(xié)議族上,是為硬件性能低下的遠程設備以及網(wǎng)絡狀況糟糕的情況下而設計的發(fā)布/訂閱型消息協(xié)議。國內很多企業(yè)都廣泛使用MQTT作為Android手機客戶端與服務器端推送消息的協(xié)議。

特點

MQTT協(xié)議是為大量計算能力有限,且工作在低帶寬、不可靠的網(wǎng)絡的遠程傳感器和控制設備通訊而設計的協(xié)議,它具有以下主要的幾項特性:

  • 使用發(fā)布/訂閱消息模式,提供一對多的消息發(fā)布,解除應用程序耦合;
  • 對負載內容屏蔽的消息傳輸;
  • 使用TCP/IP提供網(wǎng)絡連接;
  • 有三種消息發(fā)布服務質量;

至多一次:消息發(fā)布完全依賴底層 TCP/IP 網(wǎng)絡。會發(fā)生消息丟失或重復。這一級別可用于如下情況,環(huán)境傳感器數(shù)據(jù),丟失一次讀記錄無所謂,因為不久后還會有第二次發(fā)送。

至少一次:確保消息到達,但消息重復可能會發(fā)生。

只有一次:確保消息到達一次。這一級別可用于如下情況,在計費系統(tǒng)中,消息重復或丟失會導致不正確的結果。

  • 小型傳輸,開銷很?。ü潭ㄩL度的頭部是 2 字節(jié)),協(xié)議交換最小化,以降低網(wǎng)絡流量;
  • 使用Last Will和Testament特性通知有關各方客戶端異常中斷的機制。

Apache-Apollo

Apache Apollo是一個代理服務器,其是在ActiveMQ基礎上發(fā)展而來的,可以支持STOMP, AMQP, MQTT, Openwire, SSL, WebSockets 等多種協(xié)議。

原理:服務器端創(chuàng)建一個唯一訂閱號,發(fā)送者可以向這個訂閱號中發(fā)東西,然后接受者(即訂閱了這個訂閱號的人)都會收到這個訂閱號發(fā)出來的消息。以此來完成消息的推送。服務器其實是一個消息中轉站。

下載

下載地址:http://archive.apache.org/dist/activemq/activemq-apollo/

配置與啟動

  • 需要安裝JDK環(huán)境
  • 在命令行模式下進入bin,執(zhí)行apollo create mybroker d:\apache-apollo\broker,創(chuàng)建一個名為mybroker虛擬主機(Virtual Host)。需要特別注意的是,生成的目錄就是以后真正啟動程序的位置。
  • 在命令行模式下進入d:\apache-apollo\broker\bin,執(zhí)行apollo-broker run,也可以用apollo-broker-service.exe配置服務。
  • 訪問http://127.0.0.1:61680打開web管理界面。(密碼查看broker/etc/users.properties)
  • 啟動端口,看cmd輸出。

SpringBoot2的開發(fā)

添加依賴

<!-- 
  spring-boot版本 2.1.0.RELEASE
-->
<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>

自定義配置

# src/main/resources/config/mqtt.properties
##################
#  MQTT 配置
##################
# 用戶名
mqtt.username=admin
# 密碼
mqtt.password=password
# 推送信息的連接地址,如果有多個,用逗號隔開,如:tcp://127.0.0.1:61613,tcp://192.168.1.61:61613
mqtt.url=tcp://127.0.0.1:61613
##################
#  MQTT 生產(chǎn)者
##################
# 連接服務器默認客戶端ID
mqtt.producer.clientId=mqttProducer
# 默認的推送主題,實際可在調用接口時指定
mqtt.producer.defaultTopic=topic1
##################
#  MQTT 消費者
##################
# 連接服務器默認客戶端ID
mqtt.consumer.clientId=mqttConsumer
# 默認的接收主題,可以訂閱多個Topic,逗號分隔
mqtt.consumer.defaultTopic=topic1

配置MQTT發(fā)布和訂閱

import org.apache.commons.lang3.StringUtils;
import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
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;
/**
 * MQTT配置,生產(chǎn)者
 *
 * @author BBF
 */
@Configuration
public class MqttConfig {
  private static final Logger LOGGER = LoggerFactory.getLogger(MqttConfig.class);
  private static final byte[] WILL_DATA;
  static {
    WILL_DATA = "offline".getBytes();
  }
  /**
   * 訂閱的bean名稱
   */
  public static final String CHANNEL_NAME_IN = "mqttInboundChannel";
  /**
   * 發(fā)布的bean名稱
   */
  public static final String CHANNEL_NAME_OUT = "mqttOutboundChannel";
  @Value("${mqtt.username}")
  private String username;
  @Value("${mqtt.password}")
  private String password;
  @Value("${mqtt.url}")
  private String url;
  @Value("${mqtt.producer.clientId}")
  private String producerClientId;
  @Value("${mqtt.producer.defaultTopic}")
  private String producerDefaultTopic;
  @Value("${mqtt.consumer.clientId}")
  private String consumerClientId;
  @Value("${mqtt.consumer.defaultTopic}")
  private String consumerDefaultTopic;
  /**
   * 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());
    options.setServerURIs(StringUtils.split(url, ","));
    // 設置超時時間 單位為秒
    options.setConnectionTimeout(10);
    // 設置會話心跳時間 單位為秒 服務器會每隔1.5*20秒的時間向客戶端發(fā)送心跳判斷客戶端是否在線,但這個方法并沒有重連的機制
    options.setKeepAliveInterval(20);
    // 設置“遺囑”消息的話題,若客戶端與服務器之間的連接意外中斷,服務器將發(fā)布客戶端的“遺囑”消息。
    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信息通道(生產(chǎn)者)
   *
   * @return {@link org.springframework.messaging.MessageChannel}
   */
  @Bean(name = CHANNEL_NAME_OUT)
  public MessageChannel mqttOutboundChannel() {
    return new DirectChannel();
  }
  /**
   * MQTT消息處理器(生產(chǎn)者)
   *
   * @return {@link org.springframework.messaging.MessageHandler}
   */
  @Bean
  @ServiceActivator(inputChannel = CHANNEL_NAME_OUT)
  public MessageHandler mqttOutbound() {
    MqttPahoMessageHandler messageHandler = new MqttPahoMessageHandler(
        producerClientId,
        mqttClientFactory());
    messageHandler.setAsync(true);
    messageHandler.setDefaultTopic(producerDefaultTopic);
    return messageHandler;
  }
  /**
   * 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.MessageChannel}
   */
  @Bean(name = CHANNEL_NAME_IN)
  public MessageChannel mqttInboundChannel() {
    return new DirectChannel();
  }
  /**
   * MQTT消息處理器(消費者)
   *
   * @return {@link org.springframework.messaging.MessageHandler}
   */
  @Bean
  @ServiceActivator(inputChannel = CHANNEL_NAME_IN)
  public MessageHandler handler() {
    return new MessageHandler() {
      @Override
      public void handleMessage(Message<?> message) throws MessagingException {
        LOGGER.error("===================={}============", message.getPayload());
      }
    };
  }
}

消息發(fā)布器

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;
/**
 * MQTT生產(chǎn)者消息發(fā)送接口
 * <p>MessagingGateway要指定生產(chǎn)者的通道名稱</p>
 * @author BBF
 */
@Component
@MessagingGateway(defaultRequestChannel = MqttConfig.CHANNEL_NAME_OUT)
public interface IMqttSender {
  /**
   * 發(fā)送信息到MQTT服務器
   *
   * @param data 發(fā)送的文本
   */
  void sendToMqtt(String data);
  /**
   * 發(fā)送信息到MQTT服務器
   *
   * @param topic 主題
   * @param payload 消息主體
   */
  void sendToMqtt(@Header(MqttHeaders.TOPIC) String topic,
      String payload);
  /**
   * 發(fā)送信息到MQTT服務器
   *
   * @param topic 主題
   * @param qos 對消息處理的幾種機制。
 0 表示的是訂閱者沒收到消息不會再次發(fā)送,消息會丟失。

   * 1 表示的是會嘗試重試,一直到接收到消息,但這種情況可能導致訂閱者收到多次重復消息。

   * 2 多了一次去重的動作,確保訂閱者收到的消息有一次。
   * @param payload 消息主體
   */
  void sendToMqtt(@Header(MqttHeaders.TOPIC) String topic,
      @Header(MqttHeaders.QOS) int qos,
      String payload);
}

發(fā)送消息

/**
 * MQTT消息發(fā)送
 *
 * @author BBF
 */
@Controller
@RequestMapping(value = "/")
public class MqttController {
  /**
   * 注入發(fā)送MQTT的Bean
   */
  @Resource
  private IMqttSender iMqttSender;
  /**
   * 發(fā)送MQTT消息
   * @param message 消息內容
   * @return 返回
   */
  @ResponseBody
  @GetMapping(value = "/mqtt", produces ="text/html")
  public ResponseEntity<String> sendMqtt(@RequestParam(value = "msg") String message) {
    iMqttSender.sendToMqtt(message);
    return new ResponseEntity<>("OK", HttpStatus.OK);
  }
}

入口類

import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.boot.autoconfigure.jdbc.DataSourceAutoConfiguration;
import org.springframework.boot.autoconfigure.orm.jpa.HibernateJpaAutoConfiguration;
import org.springframework.context.annotation.PropertySource;
/**
 * SpringBoot 入口類
 *
 * @author BBF
 */
@SpringBootApplication(exclude = {DataSourceAutoConfiguration.class,
    HibernateJpaAutoConfiguration.class})
@PropertySource(encoding = "UTF-8", value = {"classpath:config/mqtt.properties"})
public class Application {
  public static void main(String[] args) {
    SpringApplication.run(Application.class, args);
  }
}

代碼

https://gitee.com/bbfbbf/mqtt-test

以上就是SpringBoot集成MQTT示例詳解的詳細內容,更多關于SpringBoot集成MQTT的資料請關注腳本之家其它相關文章!

相關文章

  • SpringCloudGateway網(wǎng)關處攔截并修改請求的操作方法

    SpringCloudGateway網(wǎng)關處攔截并修改請求的操作方法

    這篇文章主要介紹了SpringCloudGateway網(wǎng)關處攔截并修改請求的操作方法,本文通過示例代碼給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友參考下吧
    2023-12-12
  • Spring Cloud Feign 自定義配置(重試、攔截與錯誤碼處理) 代碼實踐

    Spring Cloud Feign 自定義配置(重試、攔截與錯誤碼處理) 代碼實踐

    這篇文章主要介紹了Spring Cloud Feign 自定義配置(重試、攔截與錯誤碼處理) 實踐,本文通過實例代碼給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2020-08-08
  • springmvc實現(xiàn)跨服務器文件上傳功能

    springmvc實現(xiàn)跨服務器文件上傳功能

    這篇文章主要為大家詳細介紹了springmvc實現(xiàn)跨服務器文件上傳功能,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2019-08-08
  • Java實現(xiàn)馬踏棋盤算法

    Java實現(xiàn)馬踏棋盤算法

    這篇文章主要為大家詳細介紹了Java實現(xiàn)馬踏棋盤算法,文中示例代碼介紹的非常詳細,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2022-02-02
  • java實現(xiàn)簡易貪吃蛇游戲

    java實現(xiàn)簡易貪吃蛇游戲

    這篇文章主要為大家詳細介紹了java實現(xiàn)簡易貪吃蛇游戲,文中示例代碼介紹的非常詳細,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2020-12-12
  • 詳談fastjson將對象格式化成json時的兩個問題

    詳談fastjson將對象格式化成json時的兩個問題

    下面小編就為大家?guī)硪黄斦刦astjson將對象格式化成json時的兩個問題。小編覺得挺不錯的,現(xiàn)在就分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2017-05-05
  • Java用?Gradle配置compile及implementation和api的區(qū)別

    Java用?Gradle配置compile及implementation和api的區(qū)別

    這篇文章主要介紹了Java用Gradle配置compile及implementation和api的區(qū)別,文章圍繞主題的相關資料展開詳細的內容介紹,具有一定的參考價值,需要的小伙伴可以參考一下
    2022-06-06
  • Java?for循環(huán)標簽跳轉到指定位置的示例詳解

    Java?for循環(huán)標簽跳轉到指定位置的示例詳解

    這篇文章主要介紹了Java?for循環(huán)標簽跳轉到指定位置,本文通過實例代碼給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2023-05-05
  • SpringBoot如何使用ApplicationContext獲取bean對象

    SpringBoot如何使用ApplicationContext獲取bean對象

    這篇文章主要介紹了SpringBoot 如何使用ApplicationContext獲取bean對象,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-11-11
  • Java中覆蓋finalize()方法實例代碼

    Java中覆蓋finalize()方法實例代碼

    這篇文章主要介紹了Java中覆蓋finalize()方法實例代碼,分享了相關代碼示例,小編覺得還是挺不錯的,具有一定借鑒價值,需要的朋友可以參考下
    2018-02-02

最新評論

天气| 临泽县| 军事| 双江| 正阳县| 云梦县| 乌鲁木齐县| 麟游县| 金坛市| 明溪县| 雷州市| 开江县| 石泉县| 麟游县| 闽清县| 旬邑县| 福安市| 岳西县| 徐州市| 上饶市| 德清县| 横峰县| 包头市| 鄂托克前旗| 恩平市| 那坡县| 东宁县| 华坪县| 许昌市| 衡南县| 浙江省| 绥滨县| 阿尔山市| 克东县| 始兴县| 仙游县| 娄底市| 泰顺县| 万年县| 潞城市| 卢龙县|