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

如何使用Spring?integration在Springboot中集成Mqtt詳解

 更新時間:2023年02月24日 09:10:52   作者:小艾咪  
MQTT是多個客戶端通過一個中央服務器傳遞信息的多對多協議,能高效地將信息分發(fā)給一個或多個訂閱者,下面這篇文章主要給大家介紹了關于如何使用Spring?integration在Springboot中集成Mqtt的相關資料,需要的朋友可以參考下

前言

接觸Mqtt最開始用的是SpingMVC 引入 Mqtt Client Jar包自己實現收發(fā)重連等操作,而現在則使用Spring Integration 進行集成。其復雜度和代碼簡潔度天差地別。所以特別寫下此文章算是對自己的一個記錄。也希望觀看到此文章的人少走一點彎路,哪怕只有一步。

當然有精力的話最好還是使用Client自己實現一遍對自己來說也是一種提高,生產環(huán)境,我建議使用Spring Integration 集成

關于Spring Intergration

介紹關于Spring Integration 的一些概念性的東西,如果只想知道怎么集成可跳過該部分

Spring Integration provides an extension of the Spring programming model to support the well known Enterprise Integration Patterns. It enables lightweight messaging within Spring-based applications and supports integration with external systems through declarative adapters. Those adapters provide a higher level of abstraction over Spring’s support for remoting, messaging, and scheduling.

大致意思是

  • Spring Integration 為許多著名的企業(yè)級集成方案提供了擴展接口
  • 它能夠在基于Sping的應用中進行輕量級消息傳遞
  • 支持聲明式適配器與發(fā)布系統進行集成(見Mqtt send消息的實現)
  • 這些適配器在 Spring對遠程調用,消息傳遞和系統調度提供了更高層次的抽象

Spring Intergration核心組件

Message(消息)

Message是對消息的包裝,在Spring 系統中傳遞的任何消息都會被包裝為Message??梢岳斫鉃槭荢pring Integration消息傳遞的基本單位

Message Channel(消息管道)

消息管道:Message 在Message Channel中進行傳遞,生產者向管道中投遞消息,消費者從管道中取出消息。Spring Integration 支持兩種消息傳遞模型,point-to-point(點對點模型),Publish-subscribe(發(fā)布訂閱模型)有多種管道類型。

本次使用點對點模型,消息管道類型DirectChannel

Message Endpoint(消息切入點)

消息切入點:消息在管道中流動那必定會有某些流入或流出的點亦或是在某個位置(即特定函數)需要對消息進行處理,過濾,格式轉換等。這些點即為Message Endpoint(實際為某些處理函數),例如消息發(fā)送,消息接收都是Message Endpoint。

Message Transformer(消息轉化器)

消息轉化器:是將消息進行特定轉換例如將一個 Object 序列化為 Json 字符串

Message Filter

消息過濾器,過濾掉特定消息。例如在管道中發(fā)送的含username 和 age 屬性的 User 對象,如果當前消息(一個User實例的包裝)的age < 18則將其過濾掉,那么處在過濾器之后的消費者將無法接收到 age < 18的User對象。當然過濾條件不僅是消息負載的屬性,也可以是消息本身的屬性。

Message Router(消息路由)

消息路由:向管道投遞消息時可由消息路由根據路由規(guī)則選擇投遞給那個管道

Splitter(分割器)

分割器:它從一個輸入管道接收一條消息并將其分割為多條消息,再將每條消息發(fā)送到不同的輸出管道上

Aggregator(聚合器)

聚合器:與分割器功能剛好相反

Service Activator

Service Activator: 它是一個用于將系統服務實例接入到消息系統的泛型切入點,該切入點必須配置輸入管道。其返回值可是消息類型也可以是一個消息處理器,當返回值為消息類型時需要指定輸出管道,即在該切入點對消息加工處理后再發(fā)送到指定的輸出管道,如果返回值為消息處理器。那么消息交由消息處理器進行處理。下文中會為Mqtt消息出站配置Service Activator并且 返回值為消息處理器

Channel Adapter(管道適配器)

管道適配器:因為外部協議有無數種,消息適配器則用于連接不同協議的外部系統。從外部系統讀入數據并對數據進行處理最終與Spring Integration 內部的消息系統適配。例如將要進行Mqtt集成,那么就需要一個Mqtt的管道適配器,事實上也確實有一個,下文中將會看到。

開始集成

依賴管理工具使用Gardle

引入spring-integration-mqtt依賴

implementation "org.springframework.integration:spring-integration-mqtt:5.4.6"

創(chuàng)建Mqtt配置類

@Configuration
public class MqttConfig {
    /**
     *  以下屬性將在配置文件中讀取
     **/
    //mqtt Broker 地址
    private String[] uris;
    //連接用戶名
    private String username;
    //連接密碼
    private String password;
    //入站Client ID
    private String inClientId;
    //出站Client ID
    private String outClientId;
    //默認訂閱主題
    private String defaultTopic;

    public void setUris(String[] uris) {
        this.uris = uris;
    }

    public void setUsername(String username) {
        this.username = username;
    }

    public void setPassword(String password) {
        this.password = password;
    }

    public void setInClientId(String inClientId) {
        this.inClientId = inClientId;
    }

    public void setOutClientId(String outClientId) {
        this.outClientId = outClientId;
    }

    public void setDefaultTopic(String defaultTopic) {
        this.defaultTopic = defaultTopic;
    }
}

這里需要注意為什么創(chuàng)建兩個client ID,Spring Integration 在集成的時候入站與出站消息處理并不使用同一個連接,所以如果clien ID相同,將會出現Mqtt反復重連現象,實為 mqtt 出入站連接交替踢對方下線。

修改配置文件 application.yml

mqtt:
  uris: tcp://ip:port
  username: user
  password: pwd
  in-client-id: ${random.value} # 隨機值,使出入站 client ID 不同
  out-client-id: ${random.value}
  default-topic: defaultTopic

在MqttConfig上使用注解 @ConfigurationProperties(prefix = "mqtt")將配置文件中屬性注入到MqttConfig中,但別忘記在啟動類上使用@EnableConfigurationProperties啟用屬性注入。

@Configuration
@ConfigurationProperties(prefix = "mqtt")
public class MqttConfig {
    .....
}

創(chuàng)建三個管道

這三個管道分別用于處理入站消息,出站消息,錯誤消息

@Bean
public MessageChannel mqttOutboundChannel(){
    return new DirectChannel();
}
@Bean
public MessageChannel mqttInboundChannel(){
    return new DirectChannel();
}
@Bean
public MessageChannel errorChannel(){
    return new DirectChannel();
}

添加 Mqtt 客戶端工廠

@Bean
public MqttPahoClientFactory mqttPahoClientFactory(){
    DefaultMqttPahoClientFactory factory = new DefaultMqttPahoClientFactory();
    MqttConnectOptions options = new MqttConnectOptions();
    options.setServerURIs(uris);
    options.setUserName(username);
    options.setPassword(password.toCharArray());
    factory.setConnectionOptions(options);
    return factory;
}

添加 Mqtt 管道適配器

@Bean
public MqttPahoMessageDrivenChannelAdapter adapter(MqttPahoClientFactory factory){
    return new MqttPahoMessageDrivenChannelAdapter(inClientId,factory,defaultTopic);
}

添加消息生產者

@Bean
public MessageProducer mqttInbound(MqttPahoMessageDrivenChannelAdapter adapter){
    adapter.setCompletionTimeout(5000);
    adapter.setConverter(new DefaultPahoMessageConverter());
    //入站投遞給入站管道
    adapter.setOutputChannel(mqttInboundChannel());
    adapter.setErrorChannel(errorChannel());
    adapter.setQos(0);
    return adapter;
}

添加出站處理器

出站處理器是一個Service Activator

@Bean
@ServiceActivator(inputChannel = "mqttOutboundChannel")
public MessageHandler mqttOutbound(MqttPahoClientFactory factory){
    MqttPahoMessageHandler handler =
            new MqttPahoMessageHandler(outClientId,factory);
    handler.setAsync(true);
    handler.setConverter(new DefaultPahoMessageConverter());
    handler.setDefaultTopic(defaultTopic);
    return handler;
}

添加消息接收器

通過前文的配置,當Mqtt 訂閱主題產生消息時會通過 MessageProducer(本例中是一個管道適配器)將消息投遞到入站管道中,所以當需要接收并處理Mqtt消息時只需要從入站管道中取出消息即可。取出消息即可使用前文的 Endpoint 本例使用Service Activator

創(chuàng)建一個接收類,自定義任意類型。重點是使用其內部的方法,將其注冊為Endpoint

@Component
public class Receiver {
    @Bean
    //使用ServiceActivator 指定接收消息的管道為 mqttInboundChannel,投遞到mqttInboundChannel管道中的消息會被該方法接收并執(zhí)行
    @ServiceActivator(inputChannel = "mqttInboundChannel")
    public MessageHandler handleMessage() {
        return message -> {
            System.out.println(message.getPayload());
        };
    }
}

至此,整個服務即可接收Mqtt broker的消息了。默認接收的主題為配置文件中指定的 "defaultTopic"。

添加消息發(fā)送器

那么消息發(fā)送器無疑是向出站管道投遞消息即可。如何實現,Spring Integration 提供了 @MessagingGateway注解,該注解提供一個defaultRequestChannel屬性用于指定出站管道。如下

@Component
@MessagingGateway(defaultRequestChannel = "mqttOutboundChannel")
public interface MqttSender {
    void sendToMqtt(String text);

    void sendWithTopic(@Header(MqttHeaders.TOPIC) String topic, String text);

    void sendWithTopicAndQos(@Header(MqttHeaders.TOPIC) String topic, @Header(MqttHeaders.QOS) Integer Qos, String text);
}

此接口內函數參數將作為消息的負載被包裝成為消息并投遞到出站管道中。同時可以看到,方法參數可以注解的方式自定義消息發(fā)送的 topic qos retain 等屬性。

至此集成完成,用戶可調用MqttSender發(fā)布消息或是沖Receiver中接收并處理消息

自定義編解碼器

前文中添加消息生產者和添加出站處理器代碼中都可以看到setConverter(new DefaultPahoMessageConverter());方法,該方法用與對消息負載進行編解碼。多數情況下我們的消息都是可編碼的。我在消息傳遞過程中使用的是json編碼。下面我將以json編碼為例,展示如何自定義編解碼器

//協議對象
public class User {
    private String username;
    private String password;

    public String getUsername() {
        return username;
    }

    public void setUsername(String username) {
        this.username = username;
    }

    public String getPassword() {
        return password;
    }

    public void setPassword(String password) {
        this.password = password;
    }
}

實現自己的編解碼器

代碼如下,詳見注釋

@Component
public class AlMingConverter implements MqttMessageConverter {
    private final  static Logger log = LoggerFactory.getLogger(AlMingConverter.class);
    private int defaultQos = 0;
    private boolean defaultRetain = false;
    ObjectMapper om = new ObjectMapper();
    //入站消息解碼
    @Override
    public Message<User> toMessage(String topic, MqttMessage mqttMessage) {
        
        User protocol = null;
        try {
            protocol = om.readValue(mqttMessage.getPayload(), User.class);
        } catch (IOException e) {
            if (e instanceof JsonProcessingException) {
                System.out.println();
                log.error("Converter only support json string");
            }
        }
        assert protocol != null;
        MessageBuilder<User> messageBuilder = MessageBuilder
                .withPayload(protocol);
        //使用withPayload初始化的消息缺少頭信息,將原消息頭信息填充進去
        messageBuilder.setHeader(MqttHeaders.ID, mqttMessage.getId())
                .setHeader(MqttHeaders.RECEIVED_QOS, mqttMessage.getQos())
                .setHeader(MqttHeaders.DUPLICATE, mqttMessage.isDuplicate())
                .setHeader(MqttHeaders.RECEIVED_RETAINED, mqttMessage.isRetained());
        if (topic != null) {
            messageBuilder.setHeader(MqttHeaders.TOPIC, topic);
        }
        return messageBuilder.build();
    }
    //出站消息編碼
    @Override
    public Object fromMessage(Message<?> message, Class<?> targetClass) {
        MqttMessage mqttMessage = new MqttMessage();
        String msg = null;
        try {
            msg = om.writeValueAsString(message.getPayload());
        } catch (JsonProcessingException e) {
            e.printStackTrace();
        }
        assert msg != null;
        mqttMessage.setPayload(msg.getBytes(StandardCharsets.UTF_8));
        //這里的 mqtt_qos ,和 mqtt_retained 由 MqttHeaders 此類得出不可以隨便取,如需其他屬性自行查找
        Integer qos = (Integer) message.getHeaders().get("mqtt_qos");
        mqttMessage.setQos(qos == null ? defaultQos : qos);
        Boolean retained = (Boolean) message.getHeaders().get("mqtt_retained");
        mqttMessage.setRetained(retained == null ? defaultRetain : retained);
        return mqttMessage;
    }
    //此方法直接拿默認編碼器的來用的,照抄即可
    @Override
    public Message<?> toMessage(Object payload, MessageHeaders headers) {
        Assert.isInstanceOf(MqttMessage.class, payload,
                () -> "This converter can only convert an 'MqttMessage'; received: " + payload.getClass().getName());
        return this.toMessage(null, (MqttMessage) payload);
    }

    public void setDefaultQos(int defaultQos) {
        this.defaultQos = defaultQos;
    }

    public void setDefaultRetain(boolean defaultRetain) {
        this.defaultRetain = defaultRetain;
    }
}

修改MqttConfig

將自定義編解碼器注入到MqttConfig中。

private final AlMingConverter alMingConverter;
@Autowired
public MqttConfig(AlMingConverter alMingConverter) {
    this.alMingConverter = alMingConverter;
}

將原來所有DefaultPahoMessageConverter實例更換為alMingConverter

...
adapter.setConverter(alMingConverter);
...
handler.setConverter(alMingConverter);

修改MqttSender

@Component
@MessagingGateway(defaultRequestChannel = "mqttOutboundChannel")
public interface MqttSender {
    void sendToMqtt(User user);

    void sendWithTopic(@Header(MqttHeaders.TOPIC) String topic, User user);

    void sendWithTopicAndQos(@Header(MqttHeaders.TOPIC) String topic, @Header(MqttHeaders.QOS) Integer Qos, User user);
}

Receiver修改

此時Reveiver的消息負載雖然編譯時為 Object ,但運行時為User可安全的進行 Object 到 User 的強轉然后進行后續(xù)操作

@Component
public class Receiver {
    @Bean
    @ServiceActivator(inputChannel = "mqttInboundChannel")
    public MessageHandler handleMessage() {
        return message -> {
            User user = (User) message.getPayload();
            System.out.println(user.getUsername()+":"+user.getPassword());
        };
    }
}

動態(tài)修改訂閱主題

前文中發(fā)布訂閱的主題都是默認主題,即在配置文件中指定的主題。生產環(huán)境肯定不止使用一個主題。那么就需要運行時動態(tài)修改主題。Spring Integration也確實為此提供了支持。上文中向Bean容器中注冊的管道適配器(MqttPahoMessageDrivenChannelAdapter)提供了addTopic removeTopic等方法可用于運行時修改主題。所以可以創(chuàng)建一個MqttService專門用于添加和刪除主題。

但注意只有Spring Integration 4.1以上版本可以這么使用

@Service
public class MqttServiceImpl implements MqttService { //MqttService是自己定義的,僅包含如下方法均已重寫
    MqttPahoMessageDrivenChannelAdapter adapter;
    @Autowired
    public MqttServiceImpl(MqttPahoMessageDrivenChannelAdapter adapter) {
        this.adapter = adapter;
    }

    @Override
    public void addTopic(String topic) {
        String[] topics = adapter.getTopic();
        if(!Arrays.asList(topics).contains(topic)){
            adapter.addTopic(topic,0);
        }
    }

    @Override
    public void removeTopic(String topic) {
        adapter.removeTopic(topic);
    }
}

增加或刪除主題時,注入該服務調用addTopic,removeTopic即可

End&附錄

最終MqttConfig完整代碼

@Configuration
@ConfigurationProperties(prefix = "mqtt")
public class MqttConfig {
    /**
     *  以下屬性將在配置文件中讀取
     **/
    //mqtt Broker 地址
    private String[] uris;
    //連接用戶名
    private String username;
    //連接密碼
    private String password;
    //入站Client ID
    private String inClientId;
    //出站Client ID
    private String outClientId;
    //默認訂閱主題
    private String defaultTopic;

    private final AlMingConverter alMingConverter;
    @Autowired
    public MqttConfig(AlMingConverter alMingConverter) {
        this.alMingConverter = alMingConverter;
    }

    @Bean
    public MqttPahoClientFactory mqttPahoClientFactory(){
        DefaultMqttPahoClientFactory factory = new DefaultMqttPahoClientFactory();
        MqttConnectOptions options = new MqttConnectOptions();
        options.setServerURIs(uris);
        options.setUserName(username);
        options.setPassword(password.toCharArray());
        factory.setConnectionOptions(options);
        return factory;
    }
    @Bean
    public MqttPahoMessageDrivenChannelAdapter adapter(MqttPahoClientFactory factory){
        return new MqttPahoMessageDrivenChannelAdapter(inClientId,factory,defaultTopic);
    }
    public void setUris(String[] uris) {
        this.uris = uris;
    }
    @Bean
    public MessageProducer mqttInbound(MqttPahoMessageDrivenChannelAdapter adapter){
        adapter.setCompletionTimeout(5000);
        adapter.setConverter(alMingConverter);
        //入站投遞的通道
        adapter.setOutputChannel(mqttInboundChannel());
        adapter.setErrorChannel(errorChannel());
        adapter.setQos(0);
        return adapter;
    }
    @Bean
    @ServiceActivator(inputChannel = "mqttOutboundChannel")
    public MessageHandler mqttOutbound(MqttPahoClientFactory factory){
        MqttPahoMessageHandler handler =
                new MqttPahoMessageHandler(outClientId,factory);
        handler.setAsync(true);
        handler.setConverter(alMingConverter);
        handler.setDefaultTopic(defaultTopic);
        return handler;
    }
    @Bean
    public MessageChannel mqttOutboundChannel(){
        return new DirectChannel();
    }
    @Bean
    public MessageChannel mqttInboundChannel(){
        return new DirectChannel();
    }
    @Bean
    public MessageChannel errorChannel(){
        return new DirectChannel();
    }
    public void setUsername(String username) {
        this.username = username;
    }

    public void setPassword(String password) {
        this.password = password;
    }

    public void setInClientId(String inClientId) {
        this.inClientId = inClientId;
    }

    public void setOutClientId(String outClientId) {
        this.outClientId = outClientId;
    }

    public void setDefaultTopic(String defaultTopic) {
        this.defaultTopic = defaultTopic;
    }
}

最后

到此這篇關于如何使用Spring integration在Springboot中集成Mqtt的文章就介紹到這了,更多相關Spring integration在Springboot集成Mqtt內容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!

相關文章

  • Java工作隊列代碼詳解

    Java工作隊列代碼詳解

    這篇文章主要介紹了Java工作隊列代碼詳解,涉及Round-robin 轉發(fā),消息應答(messageacknowledgments),消息持久化(Messagedurability)等相關內容,具有一定參考價值,需要的朋友可以了解下。
    2017-11-11
  • Java批量修改文件名的實例代碼

    Java批量修改文件名的實例代碼

    幾天前在163公開課上下了一些mp4視頻文件。發(fā)現課程名和文件名不對應,想到編個程序批量修改。先分析網頁源代碼將課程名和文件名一一對應,存儲在一個文件里,然后使用Java讀取該文件進而修改文件名。
    2013-04-04
  • 全面了解Java反射機制

    全面了解Java反射機制

    Java的反射機制在實踐中可謂無處不在,如果你已經工作幾年,還對Java的反射機制一知半解,那么這篇文章絕對值得你讀一讀。
    2020-03-03
  • ObjectMapper 如何忽略字段大小寫

    ObjectMapper 如何忽略字段大小寫

    這篇文章主要介紹了使用ObjectMapper實現忽略字段大小寫操作,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-06-06
  • Java幾個重要的關鍵字詳析

    Java幾個重要的關鍵字詳析

    這篇文章主要介紹了Java幾個重要的關鍵字詳析,文章圍繞主題展開詳細的內容介紹,具有一定的參考一下,需要的小伙伴可以參考一下,希望對你的學習有所幫助
    2022-07-07
  • Java冒泡排序法和選擇排序法的實現

    Java冒泡排序法和選擇排序法的實現

    這篇文章主要介紹了Java冒泡排序法和選擇排序法的實現,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧
    2019-09-09
  • 解決SpringBoot掃描不到公共類的實體問題

    解決SpringBoot掃描不到公共類的實體問題

    這篇文章主要介紹了解決SpringBoot掃描不到公共類的實體問題,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-08-08
  • 淺談Java中的高精度整數和高精度小數

    淺談Java中的高精度整數和高精度小數

    本篇文章主要介紹了淺談Java中的高精度整數和高精度小數,小編覺得挺不錯的,現在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2017-08-08
  • Spring中@Lazy注解的使用示例教程

    Spring中@Lazy注解的使用示例教程

    Spring在應用程序上下文啟動時去創(chuàng)建所有的單例bean對象, 而@Lazy注解可以延遲加載bean對象,即在使用時才去初始化,這篇文章主要介紹了Spring中@Lazy注解的使用,需要的朋友可以參考下
    2023-06-06
  • java面向對象設計原則之里氏替換原則示例詳解

    java面向對象設計原則之里氏替換原則示例詳解

    這篇文章主要為大家介紹了java面向對象設計原則之里氏替換原則示例詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步早日升職加薪
    2021-10-10

最新評論

富锦市| 南陵县| 合阳县| 临清市| 永德县| 越西县| 宝兴县| 绥德县| 涞水县| 满城县| 安国市| 新平| 天祝| 吐鲁番市| 阿尔山市| 赣州市| 吉木萨尔县| 榆中县| 聂拉木县| 银川市| 兴隆县| 巩留县| 遂宁市| 桂平市| 温州市| 衢州市| 洛扎县| 舞钢市| 东宁县| 房山区| 师宗县| 龙游县| 治多县| 龙口市| 连云港市| 苏尼特右旗| 泽州县| 元阳县| 明水县| 蒲江县| 巴林左旗|