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

Java連接Emqx實(shí)現(xiàn)訂閱發(fā)布消息的步驟記錄

 更新時(shí)間:2025年09月22日 10:34:31   作者:一杯冰美式_丶  
這篇文章主要介紹了Java連接Emqx實(shí)現(xiàn)訂閱發(fā)布消息的步驟記錄,EMQX是大規(guī)模分布式MQTT消息服務(wù)器,可以高效可靠連接海量物聯(lián)網(wǎng)設(shè)備,實(shí)時(shí)處理分發(fā)消息與事件流數(shù)據(jù),助力構(gòu)建關(guān)鍵業(yè)務(wù)的物聯(lián)網(wǎng)與云應(yīng)用,需要的朋友可以參考下

一:前提

安裝了Emqx開(kāi)源版、MQTTX客戶(hù)端

二:訂閱發(fā)布實(shí)現(xiàn)步驟

1.引入依賴(lài)

<!--MQTT客戶(hù)端-->
<dependency>
    <groupId>org.eclipse.paho</groupId>
    <artifactId>org.eclipse.paho.client.mqttv3</artifactId>
    <version>1.2.2</version>
</dependency>

2.編輯配置文件

mqtt:
  broker:
    uri: tcp://127.0.0.1:31883
  client:
    id: mqtt-am-client-${random.uuid}
  # 訂閱主題配置(支持多個(gè))
  inTopics:
    - topic: test/topic1
      qos: 0
    - topic: test/topic2
      qos: 1
    - topic: test/topic3
      qos: 2
  # 發(fā)布主題配置(支持多個(gè))
  outTopics:
    - topic: out/topic1
      qos: 0
  username: am
  password: LGyPtuAB4th5p
  keepAliveInterval: 60

3.讀取配置文件

package com.wtzn.web.config;

import lombok.Data;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.context.annotation.Configuration;

import java.util.List;

@Configuration
@ConfigurationProperties(prefix = "mqtt")
@Data
public class MqttProperties {
    private Broker broker;
    private Client client;
    private List<TopicConfig> inTopics;
    private List<TopicConfig> outTopics;
    private String userName;
    private String password;
    private int KeepAliveInterval;

    @Data
    public static class Broker {
        private String uri;
    }

    @Data
    public static class Client {
        private String id;
    }
    @Data
    public static class TopicConfig {
        private String topic;
        private int qos;
    }

}

4.創(chuàng)建Mqtt客戶(hù)端

package com.wtzn.web.config;

import org.eclipse.paho.client.mqttv3.MqttClient;
import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
import org.eclipse.paho.client.mqttv3.MqttException;
import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

@Configuration
public class MqttConfig {

    @Autowired
    private MqttProperties mqttProperties;

    @Bean
    public MqttClient mqttClient() throws MqttException {
        MqttClient client = new MqttClient(mqttProperties.getBroker().getUri(), mqttProperties.getClient().getId(), new MemoryPersistence());
        MqttConnectOptions options = new MqttConnectOptions();
        // 此客戶(hù)端的用戶(hù)名和密碼
        options.setUserName(mqttProperties.getUserName());
        options.setPassword(mqttProperties.getPassword().toCharArray());
        options.setCleanSession(true);
        // 設(shè)置遺囑消息
      //  options.setWill(mqttProperties.getOutTopic(), "我是mqtt-am-client,我已下線(xiàn),這是我的遺囑".getBytes(), 2, true);
        // 連接超時(shí)重試
        options.setConnectionTimeout(5000); //毫秒
        options.setKeepAliveInterval(mqttProperties.getKeepAliveInterval());
        options.setAutomaticReconnect(true);//網(wǎng)絡(luò)中斷重連
        client.connect(options);
        return client;
    }
}

5.controller層

package com.wtzn.web.controller;

import cn.dev33.satoken.annotation.SaIgnore;
import com.wtzn.common.json.utils.JsonUtils;
import com.wtzn.web.domain.bo.Payload;
import com.wtzn.web.service.MqttService;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.paho.client.mqttv3.MqttException;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.*;

import java.util.LinkedList;


@RestController
@Slf4j
@RequestMapping("/mqtt")
public class MqttController {

    @Autowired
    private MqttService mqttService;

    @SaIgnore
    @PostMapping("/mqtt")
    public void publish() {
        try {
          //  LinkedList<Payload> payloadLinkedList=new LinkedList<>();
            for(int i=1; i<=10000; i++){
                Payload payload=new Payload();
                payload.setTemperature(i);
              //  payloadLinkedList.add(payload);
                mqttService.publish("test/topic1",0,JsonUtils.toJsonString(payload));
            }

        } catch (MqttException e) {
            log.error("發(fā)布消息失敗{}", e.getMessage());
        }
        log.info("發(fā)布消息成功");
    }


}

6.service層

package com.wtzn.web.service;

import com.wtzn.common.json.utils.JsonUtils;
import com.wtzn.web.config.MqttProperties;
import com.wtzn.web.domain.bo.Payload;
import jakarta.annotation.PostConstruct;
import jakarta.annotation.PreDestroy;
import lombok.SneakyThrows;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.paho.client.mqttv3.*;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;

import java.util.Arrays;


@Service
@Slf4j
public class MqttService implements MqttCallbackExtended {

    @Autowired
    private MqttClient mqttClient;

    @Autowired
    private MqttProperties mqttProperties;
    
    @PostConstruct
    public void init() throws MqttException {
        mqttClient.setCallback(this);
 /*       mqttClient.subscribe(mqttProperties.getInTopic());
        log.info("訂閱主題{}", mqttProperties.getInTopic());
*/
        mqttProperties.getInTopics().forEach(x -> {
            try {
                mqttClient.subscribe(x.getTopic(), x.getQos());
                log.info("訂閱主題{}", x.getTopic());
            } catch (MqttException e) {
                throw new RuntimeException(e);
            }
        });

    }

    @PreDestroy
    public void destroy() throws MqttException {
        mqttClient.disconnect();
        log.info("與服務(wù)器斷開(kāi)連接");
    }

    /**
     * @description: 發(fā)送消息
     * @param: [message]
     * @return: void
     **/
    public void publish(String topic,int qos,String message) throws MqttException {
        MqttMessage mqttMessage = new MqttMessage(message.getBytes());
        mqttMessage.setQos(qos);
        mqttClient.publish(topic, mqttMessage);
        log.info("向主題【{}】發(fā)布消息:【{}】", topic, message);
    }


    /**
     * @description: 接收消息
     * @param: [topic, message]
     * @return: void
     **/
    @Override
    public void messageArrived(String topic, MqttMessage message) throws MqttException {
        Payload payload = JsonUtils.parseObject(new String(message.getPayload()), Payload.class);
        log.info("接收到來(lái)自【{}】的消息【{}】", topic, payload.getTemperature());
      /*  if (payload.getTemperature() > 37) {
            publish("發(fā)燒");
        }*/


    }


    @Override
    public void connectionLost(Throwable cause) {
        log.error("連接丟失:{}", cause.getMessage());
    }

    @SneakyThrows
    @Override
    public void deliveryComplete(IMqttDeliveryToken token) {
        if( token!=null ){
            MqttMessage message = null;
            try {
                message = token.getMessage();
            } catch (MqttException e) {
                throw new RuntimeException(e);
            }
            String topic = token.getTopics()==null ? null : Arrays.asList(token.getTopics()).toString();
            String str = message==null ? null : new String(message.getPayload());
            log.info("deliveryComplete: topic={}, message={}", topic, str);
        } else {
            log.info("deliveryComplete: null");
        }

        log.info("消息已送達(dá)");
    }

    @Override
    public void connectComplete(boolean b, String s) {

            mqttProperties.getInTopics().forEach(x -> {
                try {
                    mqttClient.subscribe(x.getTopic(), x.getQos());
                    log.info("訂閱主題{}", x.getTopic());
                } catch (MqttException e) {
                    throw new RuntimeException(e);
                }
            });
    }
}

7.dao層

package com.wtzn.web.domain.bo;

import lombok.Data;

@Data
public class Payload {
    private Integer temperature;
}

三:測(cè)試

1.PostMan直接調(diào)用測(cè)試

2、下載MQTTX客戶(hù)端進(jìn)行測(cè)試

總結(jié) 

到此這篇關(guān)于Java連接Emqx實(shí)現(xiàn)訂閱發(fā)布消息的文章就介紹到這了,更多相關(guān)Java Emqx訂閱發(fā)布消息內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Java中的WeakHashMap詳解

    Java中的WeakHashMap詳解

    這篇文章主要介紹了Java中的WeakHashMap詳解,WeakHashMap可能平時(shí)使用的頻率并不高,但是你可能聽(tīng)過(guò)WeakHashMap會(huì)進(jìn)行自動(dòng)回收吧,下面就對(duì)其原理進(jìn)行分析,需要的朋友可以參考下
    2023-09-09
  • Netty與NIO超詳細(xì)講解

    Netty與NIO超詳細(xì)講解

    Netty本質(zhì)上是一個(gè)NIO的框架,適用于服務(wù)器通訊相關(guān)的多種應(yīng)用場(chǎng)景。底層是NIO,NIO底層是Java?IO和網(wǎng)絡(luò)IO,再往下是TCP/IP協(xié)議,下面我們跟隨文章來(lái)詳細(xì)了解
    2022-08-08
  • 劍指Offer之Java算法習(xí)題精講數(shù)組與字符和等差數(shù)列

    劍指Offer之Java算法習(xí)題精講數(shù)組與字符和等差數(shù)列

    跟著思路走,之后從簡(jiǎn)單題入手,反復(fù)去看,做過(guò)之后可能會(huì)忘記,之后再做一次,記不住就反復(fù)做,反復(fù)尋求思路和規(guī)律,慢慢積累就會(huì)發(fā)現(xiàn)質(zhì)的變化
    2022-03-03
  • SpringBoot項(xiàng)目啟動(dòng)打包報(bào)錯(cuò)類(lèi)文件具有錯(cuò)誤的版本 61.0, 應(yīng)為 52.0的解決方法

    SpringBoot項(xiàng)目啟動(dòng)打包報(bào)錯(cuò)類(lèi)文件具有錯(cuò)誤的版本 61.0, 應(yīng)為 52.0的解決

    這篇文章主要給大家介紹了關(guān)于SpringBoot項(xiàng)目啟動(dòng)打包報(bào)錯(cuò)類(lèi)文件具有錯(cuò)誤的版本 61.0, 應(yīng)為 52.0的解決方法,文中有詳細(xì)的排查過(guò)程和解決方法,通過(guò)代碼介紹的非常詳細(xì),需要的朋友可以參考下
    2023-11-11
  • Spring Boot分段處理List集合多線(xiàn)程批量插入數(shù)據(jù)的解決方案

    Spring Boot分段處理List集合多線(xiàn)程批量插入數(shù)據(jù)的解決方案

    大數(shù)據(jù)量的List集合,需要把List集合中的數(shù)據(jù)批量插入數(shù)據(jù)庫(kù)中,本文給大家介紹Spring Boot分段處理List集合多線(xiàn)程批量插入數(shù)據(jù)的解決方案,感興趣的朋友跟隨小編一起看看吧
    2024-04-04
  • Java線(xiàn)程狀態(tài)變換過(guò)程代碼解析

    Java線(xiàn)程狀態(tài)變換過(guò)程代碼解析

    這篇文章主要介紹了Java線(xiàn)程狀態(tài)變換過(guò)程代碼解析,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2020-06-06
  • Java動(dòng)態(tài)線(xiàn)程池插件dynamic-tp集成過(guò)程淺析

    Java動(dòng)態(tài)線(xiàn)程池插件dynamic-tp集成過(guò)程淺析

    這篇文章主要介紹了Java動(dòng)態(tài)線(xiàn)程池插件dynamic-tp集成過(guò)程,dynamic-tp是一個(gè)輕量級(jí)的動(dòng)態(tài)線(xiàn)程池插件,它是一個(gè)基于配置中心的動(dòng)態(tài)線(xiàn)程池,線(xiàn)程池的參數(shù)可以通過(guò)配置中心配置進(jìn)行動(dòng)態(tài)的修改
    2023-03-03
  • Java項(xiàng)目中如何訪問(wèn)WEB-INF下jsp頁(yè)面

    Java項(xiàng)目中如何訪問(wèn)WEB-INF下jsp頁(yè)面

    這篇文章主要介紹了Java項(xiàng)目中如何訪問(wèn)WEB-INF下jsp頁(yè)面,文章通過(guò)示例代碼和圖文解析介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2020-08-08
  • 手把手教你如何搭建SpringBoot+Vue前后端分離

    手把手教你如何搭建SpringBoot+Vue前后端分離

    這篇文章主要介紹了手把手教你如何搭建SpringBoot+Vue前后端分離,前后端分離是目前開(kāi)發(fā)中常用的開(kāi)發(fā)模式,達(dá)成充分解耦,需要的朋友可以參考下
    2023-03-03
  • Java中Cookie和Session的那些事兒

    Java中Cookie和Session的那些事兒

    Cookie和Session都是為了保持用戶(hù)的訪問(wèn)狀態(tài),一方面為了方便業(yè)務(wù)實(shí)現(xiàn),另一方面為了簡(jiǎn)化服務(wù)端的程序設(shè)計(jì)。這篇文章主要介紹了java中cookie和session的知識(shí),需要的朋友可以參考下
    2016-09-09

最新評(píng)論

陆良县| 雷山县| 南召县| 上思县| 神木县| 从江县| 龙泉市| 奈曼旗| 洛阳市| 陇西县| 义乌市| 怀化市| 南靖县| 天门市| 定兴县| 江油市| 青铜峡市| 临湘市| 彝良县| 合川市| 漳平市| 商河县| 五莲县| 青川县| 徐州市| 万盛区| 增城市| 宜川县| 皮山县| 嘉兴市| 柞水县| 光山县| 泰来县| 铁岭县| 抚松县| 应城市| 宁城县| 阳江市| 洪泽县| 开江县| 花垣县|