SpringBoot 集成 Emqx 发布/订阅数据_springboot 整合 emqx发布 订阅-程序员宅基地

技术标签: emqx  spring boot  SpringBoot  

        最近项目中用到Emqx发布/订阅数据,特此记录便于日后查阅。

        ThingsboardEmqxTransportApplication

/**
 * Copyright  2016-2023 The Thingsboard Authors
 * <p>
 * Licensed under the Apache License, Version 2.0 (the "License");
 * you may not use this file except in compliance with the License.
 * You may obtain a copy of the License at
 * <p>
 * http://www.apache.org/licenses/LICENSE-2.0
 * <p>
 * Unless required by applicable law or agreed to in writing, software
 * distributed under the License is distributed on an "AS IS" BASIS,
 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
 * See the License for the specific language governing permissions and
 * limitations under the License.
 */
package org.thingsboard.server.emqx;

import org.springframework.boot.SpringApplication;
import org.springframework.boot.SpringBootConfiguration;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.context.annotation.ComponentScan;
import org.springframework.integration.config.EnableIntegration;

import java.util.Arrays;

@SpringBootConfiguration
@EnableConfigurationProperties
@EnableIntegration
@ComponentScan({"org.thingsboard.server.transport.emqx"})
public class ThingsboardEmqxTransportApplication {

    private static final String SPRING_CONFIG_NAME_KEY = "--spring.config.name";
    private static final String DEFAULT_SPRING_CONFIG_PARAM = SPRING_CONFIG_NAME_KEY + "=" + "tb-emqx-transport";

    public static void main(String[] args) {
        SpringApplication.run(ThingsboardEmqxTransportApplication.class, updateArguments(args));
    }

    private static String[] updateArguments(String[] args) {
        if (Arrays.stream(args).noneMatch(arg -> arg.startsWith(SPRING_CONFIG_NAME_KEY))) {
            String[] modifiedArgs = new String[args.length + 1];
            System.arraycopy(args, 0, modifiedArgs, 0, args.length);
            modifiedArgs[args.length] = DEFAULT_SPRING_CONFIG_PARAM;
            return modifiedArgs;
        }
        return args;
    }
}

        GMqttPahoMessageDrivenChannelAdapter

package org.thingsboard.server.transport.emqx.adpter;

import lombok.extern.slf4j.Slf4j;
import org.springframework.integration.mqtt.core.MqttPahoClientFactory;
import org.springframework.integration.mqtt.inbound.MqttPahoMessageDrivenChannelAdapter;

/**
 * @author zhangzhixiang on 2023/4/12
 */
@Slf4j
public class GMqttPahoMessageDrivenChannelAdapter extends MqttPahoMessageDrivenChannelAdapter {

    public GMqttPahoMessageDrivenChannelAdapter(String url, String clientId, MqttPahoClientFactory clientFactory,
                                                String... topic) {
        super(url, clientId, clientFactory, topic);
    }

    /**
     * Fix Bug.
     *
     * <p> 死锁描述:
     * Found one Java-level deadlock:
     * =============================
     * "MQTT Rec: iot-shadow-restapi_sub_hqxzgpcy":
     * waiting for ownable synchronizer 0x00000000d73e9d70, (a java.util.concurrent.locks.ReentrantLock$NonfairSync),
     * which is held by "main"
     * "main":
     * waiting to lock monitor 0x00007f5840008bf8 (object 0x00000000d73d2480, a org.springframework.integration
     * .mqtt.inbound.MqttPahoMessageDrivenChannelAdapter),
     * which is held by "MQTT Rec: iot-shadow-restapi_sub_hqxzgpcy"
     *
     * <p> 原因分析:
     * main主线程
     * AbstractEndpoint.start()获取到了ReentrantLock锁
     * MqttPahoMessageDrivenChannelAdapter.scheduleReconnect()但是需要MqttPahoMessageDrivenChannelAdapter对象锁
     * MQTT Rec线程
     * 获取到了MqttPahoMessageDrivenChannelAdapter对象锁,但是需要ReentrantLock锁
     *
     * @param cause
     */
    @Override
    public void connectionLost(Throwable cause) {
        try {
            this.lifecycleLock.lock();
        } catch (Exception e) {
            log.error("Stack Trace: {}", e);
        } finally {
            this.lifecycleLock.unlock();
        }
        super.connectionLost(cause);
    }
}

          MqttConfig

package org.thingsboard.server.transport.emqx.config;

import lombok.Data;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.RandomStringUtils;
import org.apache.commons.lang3.StringUtils;
import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.context.properties.ConfigurationProperties;
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 org.thingsboard.server.transport.emqx.DefaultMqttSubMessageHandler;
import org.thingsboard.server.transport.emqx.MqttSubMessageHandler;
import org.thingsboard.server.transport.emqx.adpter.GMqttPahoMessageDrivenChannelAdapter;
import org.thingsboard.server.transport.emqx.constant.Qos;

import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
import javax.net.ssl.SSLSocketFactory;

/**
 * @author zhangzhixiang on 2023/4/12
 */
@Data
@Slf4j
@Configuration
@ConfigurationProperties(prefix = "emqx")
@IntegrationComponentScan
public class MqttConfig {

    private String url;
    private String username;
    private String password;
    private int timeout;
    private int keepalive;
    private String enabled;
    private String subClientId;
    private String pubClientId;
    private MqttSubMessageHandler messageHandler;
    private MqttPahoMessageDrivenChannelAdapter adapter;
    private boolean dataVerifyEnabled;
    private boolean sslVerifyEnabled;
    private String caCertFile;
    private String clientCertFile;
    private String clientKeyFile;
    private String[] defaultTopics;

    ///
    // mqtt subscribe
    ///
    @PostConstruct
    public void init() {
        inbound();
        addTopic(defaultTopics);
        setMessageHandler(new DefaultMqttSubMessageHandler());
        log.info("EMQX transport started!");
    }

    @PreDestroy
    public void shutdown() {
        log.info("Stopping EMQX transport!");
        removeTopic(defaultTopics);
        log.info("EMQX transport stopped!");
    }

    @Bean
    public MessageChannel mqttInputChannel() {
        return new DirectChannel();
    }

    @Bean
    public MessageProducer inbound() {
        this.subClientId = createMqttClientId(false);
        adapter = new GMqttPahoMessageDrivenChannelAdapter(url, subClientId, mqttClientFactory());
        adapter.setCompletionTimeout(5000);
        adapter.setConverter(new DefaultPahoMessageConverter());
        adapter.setQos(Qos.DEFAULT.getValue());
        adapter.setOutputChannel(mqttInputChannel());
        log.info("Success to initialize Mqtt channel adapter for subscribe, " +
                "url=[{}],subClientId=[{}],sslVerifyEnabled=[{}]", url, subClientId, sslVerifyEnabled);
        return adapter;
    }

    @Bean
    @ServiceActivator(inputChannel = "mqttInputChannel")
    public MessageHandler handler() {
        return new MessageHandler() {
            @Override
            public void handleMessage(Message<?> message) throws MessagingException {
                messageHandler.receiveMessage(message);
            }
        };
    }

    @Bean
    public MqttPahoClientFactory mqttClientFactory() {
        DefaultMqttPahoClientFactory factory = new DefaultMqttPahoClientFactory();
        MqttConnectOptions options = new MqttConnectOptions();
        options.setCleanSession(true);
        options.setUserName(username);
        options.setPassword(password.toCharArray());
        options.setConnectionTimeout(timeout);
        options.setKeepAliveInterval(keepalive);
        // automaticReconnect 为 true 表示断线自动重连,但仅仅只是重新连接,并不订阅主题;在 connectComplete 回调函数重新订阅
        options.setAutomaticReconnect(true);
        options.setServerURIs(new String[]{url});
        if (sslVerifyEnabled) {
            try {
                SSLSocketFactory sslSocketFactory = SSLFellow.createSSLSocketFactory(
                        caCertFile, clientCertFile, clientKeyFile);
                options.setSocketFactory(sslSocketFactory);
            } catch (Exception e) {
                log.error("Stack Trace: {}", e);
            }
        }
        factory.setConnectionOptions(options);
        return factory;
    }

    ///
    // mqtt publish
    ///
    @Bean
    @ServiceActivator(inputChannel = "mqttOutboundChannel")
    public MessageHandler mqttOutbound() {
        this.pubClientId = createMqttClientId(true);
        MqttPahoMessageHandler messageHandler =
                new MqttPahoMessageHandler(pubClientId, mqttClientFactory());
        messageHandler.setAsync(true);
        messageHandler.setDefaultQos(Qos.DEFAULT.getValue());
        log.info("Success to initialize Mqtt channel adapter for publish, " +
                "url=[{}], pubClientId=[{}]", url, pubClientId);
        return messageHandler;
    }

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

    ///
    // public functions
    ///
    public void addTopic(String... topics) {
        this.adapter.addTopic(topics);
    }

    public void removeTopic(String... topics) {
        this.adapter.removeTopic(topics);
    }

    ///
    // private functions
    ///
    private String createMqttClientId(boolean publish) {
        String moduleName = Module.getName();
        if (StringUtils.isBlank(moduleName)) {
            moduleName = "iot_mqtt";
        }
        moduleName += publish ? "_pub_" : "_sub_";
        return moduleName + RandomStringUtils.randomAlphanumeric(8).toLowerCase();
    }
}

        MqttSender

package org.thingsboard.server.transport.emqx.config;

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;

/**
 * @author zhangzhixiang on 2023/4/12
 */
@Component
@MessagingGateway(defaultRequestChannel = "mqttOutboundChannel")
public interface MqttSender {

    void sendToMqtt(String data);

    void sendToMqtt(@Header(MqttHeaders.TOPIC) String topic, String payload);

    void sendToMqtt(@Header(MqttHeaders.TOPIC) String topic, @Header(MqttHeaders.QOS) int qos, String payload);
}

        到此SpringBoot 集成 Emqx 发布/订阅数据介绍完成。

版权声明:本文为博主原创文章,遵循 CC 4.0 BY-SA 版权协议,转载请附上原文出处链接和本声明。
本文链接:https://blog.csdn.net/qq_19734597/article/details/130110130

智能推荐

使用nginx解决浏览器跨域问题_nginx不停的xhr-程序员宅基地

文章浏览阅读1k次。通过使用ajax方法跨域请求是浏览器所不允许的,浏览器出于安全考虑是禁止的。警告信息如下:不过jQuery对跨域问题也有解决方案,使用jsonp的方式解决,方法如下:$.ajax({ async:false, url: 'http://www.mysite.com/demo.do', // 跨域URL ty..._nginx不停的xhr

在 Oracle 中配置 extproc 以访问 ST_Geometry-程序员宅基地

文章浏览阅读2k次。关于在 Oracle 中配置 extproc 以访问 ST_Geometry,也就是我们所说的 使用空间SQL 的方法,官方文档链接如下。http://desktop.arcgis.com/zh-cn/arcmap/latest/manage-data/gdbs-in-oracle/configure-oracle-extproc.htm其实简单总结一下,主要就分为以下几个步骤。..._extproc

Linux C++ gbk转为utf-8_linux c++ gbk->utf8-程序员宅基地

文章浏览阅读1.5w次。linux下没有上面的两个函数,需要使用函数 mbstowcs和wcstombsmbstowcs将多字节编码转换为宽字节编码wcstombs将宽字节编码转换为多字节编码这两个函数,转换过程中受到系统编码类型的影响,需要通过设置来设定转换前和转换后的编码类型。通过函数setlocale进行系统编码的设置。linux下输入命名locale -a查看系统支持的编码_linux c++ gbk->utf8

IMP-00009: 导出文件异常结束-程序员宅基地

文章浏览阅读750次。今天准备从生产库向测试库进行数据导入,结果在imp导入的时候遇到“ IMP-00009:导出文件异常结束” 错误,google一下,发现可能有如下原因导致imp的数据太大,没有写buffer和commit两个数据库字符集不同从低版本exp的dmp文件,向高版本imp导出的dmp文件出错传输dmp文件时,文件损坏解决办法:imp时指定..._imp-00009导出文件异常结束

python程序员需要深入掌握的技能_Python用数据说明程序员需要掌握的技能-程序员宅基地

文章浏览阅读143次。当下是一个大数据的时代,各个行业都离不开数据的支持。因此,网络爬虫就应运而生。网络爬虫当下最为火热的是Python,Python开发爬虫相对简单,而且功能库相当完善,力压众多开发语言。本次教程我们爬取前程无忧的招聘信息来分析Python程序员需要掌握那些编程技术。首先在谷歌浏览器打开前程无忧的首页,按F12打开浏览器的开发者工具。浏览器开发者工具是用于捕捉网站的请求信息,通过分析请求信息可以了解请..._初级python程序员能力要求

Spring @Service生成bean名称的规则(当类的名字是以两个或以上的大写字母开头的话,bean的名字会与类名保持一致)_@service beanname-程序员宅基地

文章浏览阅读7.6k次,点赞2次,收藏6次。@Service标注的bean,类名:ABDemoService查看源码后发现,原来是经过一个特殊处理:当类的名字是以两个或以上的大写字母开头的话,bean的名字会与类名保持一致public class AnnotationBeanNameGenerator implements BeanNameGenerator { private static final String C..._@service beanname

随便推点

二叉树的各种创建方法_二叉树的建立-程序员宅基地

文章浏览阅读6.9w次,点赞73次,收藏463次。1.前序创建#include&lt;stdio.h&gt;#include&lt;string.h&gt;#include&lt;stdlib.h&gt;#include&lt;malloc.h&gt;#include&lt;iostream&gt;#include&lt;stack&gt;#include&lt;queue&gt;using namespace std;typed_二叉树的建立

解决asp.net导出excel时中文文件名乱码_asp.net utf8 导出中文字符乱码-程序员宅基地

文章浏览阅读7.1k次。在Asp.net上使用Excel导出功能,如果文件名出现中文,便会以乱码视之。 解决方法: fileName = HttpUtility.UrlEncode(fileName, System.Text.Encoding.UTF8);_asp.net utf8 导出中文字符乱码

笔记-编译原理-实验一-词法分析器设计_对pl/0作以下修改扩充。增加单词-程序员宅基地

文章浏览阅读2.1k次,点赞4次,收藏23次。第一次实验 词法分析实验报告设计思想词法分析的主要任务是根据文法的词汇表以及对应约定的编码进行一定的识别,找出文件中所有的合法的单词,并给出一定的信息作为最后的结果,用于后续语法分析程序的使用;本实验针对 PL/0 语言 的文法、词汇表编写一个词法分析程序,对于每个单词根据词汇表输出: (单词种类, 单词的值) 二元对。词汇表:种别编码单词符号助记符0beginb..._对pl/0作以下修改扩充。增加单词

android adb shell 权限,android adb shell权限被拒绝-程序员宅基地

文章浏览阅读773次。我在使用adb.exe时遇到了麻烦.我想使用与bash相同的adb.exe shell提示符,所以我决定更改默认的bash二进制文件(当然二进制文件是交叉编译的,一切都很完美)更改bash二进制文件遵循以下顺序> adb remount> adb push bash / system / bin /> adb shell> cd / system / bin> chm..._adb shell mv 权限

投影仪-相机标定_相机-投影仪标定-程序员宅基地

文章浏览阅读6.8k次,点赞12次,收藏125次。1. 单目相机标定引言相机标定已经研究多年,标定的算法可以分为基于摄影测量的标定和自标定。其中,应用最为广泛的还是张正友标定法。这是一种简单灵活、高鲁棒性、低成本的相机标定算法。仅需要一台相机和一块平面标定板构建相机标定系统,在标定过程中,相机拍摄多个角度下(至少两个角度,推荐10~20个角度)的标定板图像(相机和标定板都可以移动),即可对相机的内外参数进行标定。下面介绍张氏标定法(以下也这么称呼)的原理。原理相机模型和单应矩阵相机标定,就是对相机的内外参数进行计算的过程,从而得到物体到图像的投影_相机-投影仪标定

Wayland架构、渲染、硬件支持-程序员宅基地

文章浏览阅读2.2k次。文章目录Wayland 架构Wayland 渲染Wayland的 硬件支持简 述: 翻译一篇关于和 wayland 有关的技术文章, 其英文标题为Wayland Architecture .Wayland 架构若是想要更好的理解 Wayland 架构及其与 X (X11 or X Window System) 结构;一种很好的方法是将事件从输入设备就开始跟踪, 查看期间所有的屏幕上出现的变化。这就是我们现在对 X 的理解。 内核是从一个输入设备中获取一个事件,并通过 evdev 输入_wayland

推荐文章

热门文章

相关标签