一分钟在SpringBoot项目中使用EMQ

先展示最终的结果:

生产者端:

@RestController
@RequiredArgsConstructor
public class TestController {private final MqttProducer mqttProducer;@GetMapping("/test")public String test() {User build = User.builder().age(100).sex(1).address("世界潍坊渤海之眼").build();// 延时发布mqttProducer.send("$delayed/10/cookie", 2, JSON.toJSONString(build));return "ok";}}

消费者端

/*** @author : Cookie* date : 2024/1/30*/
@Component
@Topic("cookie")
public class TestTopicHandler implements MsgHandler {@Overridepublic void process(String jsonMsg) {User user = JSON.parseObject(jsonMsg, User.class);System.out.println(user);}}

控制台结果:

在这里插入图片描述

具体解释在之前的笔记中, 需要的话可以查看 EMQ的介绍及整合SpringBoot的使用-CSDN博客


OK, 下面我们就开始实现上面的逻辑, 你要做的就是把 1-9 复制到项目即可

1. 依赖导入
<dependency><groupId>org.eclipse.paho</groupId><artifactId>org.eclipse.paho.client.mqttv3</artifactId><version>1.2.5</version>
</dependency>
2. yml 配置
# 顶格
mqtt:client:username: adminpassword: publicserverURI: tcp://192.168.200.128:1883clientId: monitor.task.${random.int[10000,99999]} # 注意: emq的客户端id 不能重复keepAliveInterval: 10  #连接保持检查周期  秒connectionTimeout: 30 #连接超时时间  秒producer:defaultQos: 2defaultRetained: falsedefaultTopic: topic/test1consumer:consumerTopics: $queue/cookie/#, $share/group1/yfs1024  #不带群组的共享订阅    多个主题逗号隔开# $queue/cookie/## 以$queue开头,不带群组的共享订阅   多个客户端只能有一个消费者消费# $share/group1/yfs1024# 以$share开头,群组的共享订阅 多个客户端订阅# 如果在一个组 只能有一个消费者消费# 如果不在一个组 都可以消费
3. 属性配置
@Data
@Configuration
@ConfigurationProperties(prefix = "mqtt.client")
public class MqttConfigProperties {private int defaultProducerQos;private boolean defaultRetained;private String defaultTopic;private String username;private String password;private String serverURI;private String clientId;private int keepAliveInterval;private int connectionTimeout;}
4. 定义主题注解
@Documented
@Retention(RetentionPolicy.RUNTIME)
@Target({ElementType.TYPE})
public @interface Topic {String value();
}
5.Mqtt配置类
@Data
@Slf4j
@Configuration
@RequiredArgsConstructor
public class MqttConfig {private final MqttConfigProperties configProperties;private final MqttCallback mqttCallback;@Beanpublic MqttClient mqttClient() {try {MqttClient client = new MqttClient(configProperties.getServerURI(), configProperties.getClientId(), mqttClientPersistence());client.setManualAcks(true); //设置手动消息接收确认mqttCallback.setMqttClient(client);client.setCallback(mqttCallback);client.connect(mqttConnectOptions());return client;} catch (MqttException e) {log.error("emq connect error", e);return null;}}@Beanpublic MqttConnectOptions mqttConnectOptions() {MqttConnectOptions options = new MqttConnectOptions();options.setUserName(configProperties.getUsername());options.setPassword(configProperties.getPassword().toCharArray());options.setAutomaticReconnect(true);//是否自动重新连接options.setCleanSession(true);//是否清除之前的连接信息options.setConnectionTimeout(configProperties.getConnectionTimeout());//连接超时时间options.setKeepAliveInterval(configProperties.getKeepAliveInterval());//心跳options.setMqttVersion(MqttConnectOptions.MQTT_VERSION_3_1_1);//设置mqtt版本return options;}public MqttClientPersistence mqttClientPersistence() {return new MemoryPersistence();}}
6. 定义消息处理接口
/*** 消息处理接口*/
public interface MsgHandler {void process(String jsonMsg) throws IOException;}
7.定义消息上下文
/*** 消息处理上下文, 通过主题拿到topic*/
public interface MsgHandlerContext{MsgHandler getMsgHandler(String topic);}
8. 定义回调类
@Component
@Slf4j
public class MqttCallback implements MqttCallbackExtended {// 需要订阅的topic配置@Value("${mqtt.consumer.consumerTopics}")private List<String> consumerTopics;@Autowiredprivate MsgHandlerContext msgHandlerContext;@Overridepublic void connectionLost(Throwable throwable) {log.error("emq error.", throwable);}@Overridepublic void messageArrived(String topic, MqttMessage message) throws Exception {log.info("topic:" + topic + "  message:" + new String(message.getPayload()));//处理消息String msgContent = new String(message.getPayload());log.info("接收到消息:" + msgContent);try {// 根据主题名称 获取 该主题对应的处理器对象// 多态 父类的引用指向子类的对象MsgHandler msgHandler = msgHandlerContext.getMsgHandler(topic);if (msgHandler == null) {return;}msgHandler.process(msgContent); //执行} catch (IOException e) {log.error("process msg error,msg is: " + msgContent, e);}//处理成功后确认消息mqttClient.messageArrivedComplete(message.getId(), message.getQos());}@Overridepublic void deliveryComplete(IMqttDeliveryToken iMqttDeliveryToken) {log.info("deliveryComplete---------" + iMqttDeliveryToken.isComplete());}@Overridepublic void connectComplete(boolean b, String s) {log.info("连接成功");//和EMQ连接成功后根据配置自动订阅topicif (consumerTopics != null && consumerTopics.size() > 0) {// 循环遍历当前项目中配置的所有的主题.consumerTopics.forEach(t -> {try {log.info(">>>>>>>>>>>>>>subscribe topic:" + t);// 订阅当前集群中所有的主题 消息服务质量 2 -> 至少收到一个mqttClient.subscribe(t, 2);} catch (MqttException e) {log.error("emq connect error", e);}});}}private MqttClient mqttClient;// 在配置类中调用传入连接public void setMqttClient(MqttClient mqttClient) {this.mqttClient = mqttClient;}
}
8. 消息处理类加载器

作用: 将Topic跟处理类对应 通过 handlerMap

/*** 消息处理类加载器* 作用:* 1. 因为实现了Spring 的 ApplicationContextAware 接口, 项目启动后就会运行实现的方法* 2. 获取MsgHandler接口的所有的实现类* 3. 将实现类上的Topic注解的值,作为handlerMap的键,实现类(处理器)作为对应的值*/
@Component
public class MsgHandlerContextImp implements ApplicationContextAware, MsgHandlerContext {private Map<String, MsgHandler> handlerMap = Maps.newHashMap();@Overridepublic void setApplicationContext(ApplicationContext applicationContext) throws BeansException {// 从spring容器中获取 <所有> 实现了MsgHandler接口的对象// key 默认类名首字母小写 value 当前对象Map<String, MsgHandler> map = applicationContext.getBeansOfType(MsgHandler.class);map.values().forEach(obj -> {// 通过反射拿到注解中的值  即 当前类处理的 topicString topic =  obj.getClass().getAnnotation(Topic.class).value();// 将主题和当前主题的处理类建立映射handlerMap.put(topic,obj);});}@Overridepublic MsgHandler getMsgHandler(String topic) {return handlerMap.get(topic);}
}
9. 封装消息生产者
@Slf4j
@Component
public class MqttProducer {// @Value() 读取配置 当然也可以批量读取配置,这里就一个一个了@Value("${mqtt.producer.defaultQos}")private int defaultProducerQos;@Value("${mqtt.producer.defaultRetained}")private boolean defaultRetained;@Value("${mqtt.producer.defaultTopic}")private String defaultTopic;@Autowiredprivate MqttClient mqttClient;public void send(String payload) {this.send(defaultTopic, payload);}public void send(String topic, String payload) {this.send(topic, defaultProducerQos, payload);}public void send(String topic, int qos, String payload) {this.send(topic, qos, defaultRetained, payload);}public void send(String topic, int qos, boolean retained, String payload) {try {mqttClient.publish(topic, payload.getBytes(), qos, retained);} catch (MqttException e) {log.error("publish msg error.",e);}}public <T> void send(String topic, int qos, T msg) throws JsonProcessingException {String payload = JsonUtil.serialize(msg);this.send(topic,qos,payload);}
}

最终的实现的结果

  • 生产者端: 在需要发送消息的地方注入 MqttProducer 发送消息
  • 消费者端: 在需要处理对应主题的类上 实现 MsgHandler接口
代码示例
生产者端
@RestController
@RequiredArgsConstructor
public class TestController {private final MqttProducer mqttProducer;@GetMapping("/test")public String test() {User build = User.builder().age(100).sex(1).address("世界潍坊渤海之眼").build();// 延时发布mqttProducer.send("$delayed/10/cookie", 2, JSON.toJSONString(build));return "ok";}
}
消费者端
@Component
@Topic("cookie")
public class TestTopicHandler implements MsgHandler {@Overridepublic void process(String jsonMsg) {User user = JSON.parseObject(jsonMsg, User.class);System.out.println(user);}}

控制台结果展示:
在这里插入图片描述

本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若转载,请注明出处:http://www.mzph.cn/news/659158.shtml

如若内容造成侵权/违法违规/事实不符,请联系多彩编程网进行投诉反馈email:809451989@qq.com,一经查实,立即删除!

相关文章

中国建设银行,这年终奖噶噶高!!!!(含算法原题)

国企年终 今天刷到一个近期帖子:「中国建设银行&#xff0c;这年终奖噶噶高!!!!」 先撇去具体内容不看&#xff0c;能在自然年的 月初&#xff0c;就把去年的奖金发了的企业&#xff0c;首先值得一个点赞。 再细看内容&#xff0c;年终奖是一个 字头的 位数。 由于国企通常没…

burp靶场--xss下篇【16-30】

burp靶场–xss下篇【16-30】 https://portswigger.net/web-security/all-labs#cross-site-scripting 实验16&#xff1a;允许使用一些 SVG 标记的反射型 XSS ### 实验要求&#xff1a; 该实验室有一个简单的反射型 XSS漏洞。该网站阻止了常见标签&#xff0c;但错过了一些 S…

Excel没有内置统计字数功能,但可以用一些变通的方法

是否需要计算Excel工作簿中某个单元格或单元格范围内的单词数? 出于多种原因,你可能需要计算文本数据中的字数。也许你有逗号分隔的列表,需要计算每个列表中的项目数。 不幸的是,Excel没有内置的单词计数方法。但是有一些聪明的方法可以得到你需要的结果。 这篇文章将向…

三步实现 Sentinel-Nacos 持久化

一、背景 版本&#xff1a;【Sentinel-1.8.6】 模式&#xff1a;【Push 模式】 参照官网介绍&#xff1a;生产环境下使用Sentinel &#xff0c;规则管理及推送模式有以下3种模式&#xff1a; 比较之后&#xff0c;目前微服务都使用了各种各样的配置中心&#xff0c;故采用Pus…

springboot综合案例(一)

文章目录 前言项目开发流程需求分析库表设计编码环节环境搭建mybatis的配置jsp模版引擎的配置日志的配置基本项目工程的配置 功能实现用户注册实现验证码功能实现用户注册 用户登录功能员工列表实现员工信息增删查改员工增加信息员工修改信息删除员工信息 前言 我具体用一个小…

【springboot图书个性化推荐系统】

前言 &#x1f31e;博主介绍&#xff1a;✌全网粉丝15W,CSDN特邀作者、211毕业、高级全栈开发程序员、大厂多年工作经验、码云/掘金/华为云/阿里云/InfoQ/StackOverflow/github等平台优质作者、专注于Java、小程序技术领域和毕业项目实战&#xff0c;以及程序定制化开发、全栈…

WindTerm 安装使用教程

一、WindTerm 功能介绍 WindTerm 是一款 Github 上开源的 SSH 终端工具&#xff0c;它是完全可以比肩 MobaXterm 工具的。其支持的系统及功能如下&#xff1a; 功能支持&#xff1a; SSHTelnetShellTCPSerialSFTPCmdPowerShellGit 二、WindTerm 官网下载 有两种获取方法&…

SpringBoot集成MongoDB(3)|(MongoTemplate的List操作)

SpringBoot集成MongoDB&#xff08;3&#xff09;|&#xff08;MongoTemplate的List操作&#xff09; 文章目录 SpringBoot集成MongoDB&#xff08;3&#xff09;|&#xff08;MongoTemplate的List操作&#xff09;[TOC] 前言一、场景说明一、向数组字段添加元素二、从数组中删…

机器学习 低代码 ML:PyCaret 的使用

✅作者简介&#xff1a;人工智能专业本科在读&#xff0c;喜欢计算机与编程&#xff0c;写博客记录自己的学习历程。 &#x1f34e;个人主页&#xff1a;小嗷犬的个人主页 &#x1f34a;个人网站&#xff1a;小嗷犬的技术小站 &#x1f96d;个人信条&#xff1a;为天地立心&…

VirtualBox中Ubuntu硬盘扩容

1.选中要扩容的虚拟机点击属性按钮&#xff0c;选择存储后点击控制器&#xff1a;STAT右边的 按钮 2.创建虚拟硬盘 在弹出框中选择创建按钮&#xff0c;选择VDI后点击下一步按钮 选择动态分配后点击下一步按钮 3.设置文件位置和大小 选择要保存的虚拟硬盘文件路径&#xff0c…

会计试算平衡

目录 一. 试算平衡的意义二. 试算平衡的原理和内容三. 试算平衡表 \quad 一. 试算平衡的意义 \quad ①验证错误 ②便于编制会计报表 试算表根据各分类账借贷余额汇总编制而成&#xff0c;依据试算表编制会计报表将比直接依据分类账来编制更为方便,拥有大量分类账的企业尤为便捷…

basic CNN

文章目录 回顾卷积神经网络卷积卷积核卷积过程卷积后图像尺寸计算公式&#xff1a;代码 padding代码 Stride代码 MaxPooling代码 一个简单的卷积神经网络用卷积神经网络来对MINIST数据集进行分类如何使用GPU代码 练习 回顾 下面这种由线形层构成的网络是全连接网络。 对于图像…

分治 (地毯填补问题)

地毯填补问题 题目描述 相传在一个古老的阿拉伯国家里&#xff0c;有一座宫殿。宫殿里有个四四方方的格子迷宫&#xff0c;国王选择驸马的方法非常特殊&#xff0c;也非常简单&#xff1a;公主就站在其中一个方格子上&#xff0c;只要谁能用地毯将除公主站立的地方外的所有地…

万户 ezOFFICE DocumentEdit_unite.jsp SQL注入漏洞复现

0x01 产品简介 万户OA ezoffice是万户网络协同办公产品多年来一直将主要精力致力于中高端市场的一款OA协同办公软件产品,统一的基础管理平台,实现用户数据统一管理、权限统一分配、身份统一认证。统一规划门户网站群和协同办公平台,将外网信息维护、客户服务、互动交流和日…

Python下载安装与环境配置

本文将指导您完成Python的下载、安装以及环境配置过程&#xff0c;确保您在编写和运行Python代码时能够获得最佳体验。我们将提供详细的步骤和代码示例&#xff0c;帮助您顺利完成设置。 一、Python下载与安装 访问Python官网&#xff1a;首先&#xff0c;您需要访问Python的官…

Pycharm 关闭/退出烦人的Pytest模式

Pycharm 遇到&#xff1a;Run Python tests in ***.py &#xff0c;但很多时候我们并不需要&#xff0c;真心烦人&#xff01; 如何解决: 1 打开File-Settings &#xff08;图片是新版界面&#xff0c;旧版同样操作&#xff09; 2 Tools 中的Python Integrated Tools 在Tes…

LeetCode —— 137. 只出现一次的数字 II

&#x1f636;‍&#x1f32b;️&#x1f636;‍&#x1f32b;️&#x1f636;‍&#x1f32b;️&#x1f636;‍&#x1f32b;️Take your time ! &#x1f636;‍&#x1f32b;️&#x1f636;‍&#x1f32b;️&#x1f636;‍&#x1f32b;️&#x1f636;‍&#x1f32b;️…

第17次修改了可删除可持久保存的前端html备忘录:增加年月日星期,增加倒计时,更改保存区名称可以多个备忘录保存不一样的信息,匹配背景主题:现代深色

第17次修改了可删除可持久保存的前端html备忘录&#xff1a;增加年月日星期&#xff0c;增加倒计时&#xff0c;更改保存区名称可以多个备忘录保存不一样的信息&#xff0c;匹配背景主题&#xff1a;现代深色 备忘录代码&#xff1a; <!DOCTYPE html> <html lang&quo…

“死“社群先不要扔,想办法激活一下,隔壁的运营都馋哭了

私域运营已成为当下很多企业寻求增长的标配。在这过程中&#xff0c;社群运营就是极为重要的一个环节。过去我们为了流量&#xff0c;疯狂建群拉人。但建社群容易活跃难&#xff0c;活跃一段时间后&#xff0c;社群会越来越安静。 不仅如此&#xff0c;群主和管理员也渐渐疏于…

c++ 字符串切分split

c 字符串切分split 的举例实现 一共给出了四种方式 1、 strtok 2、 stringstream 3、 字符串查找 4、 基于封装的方式&#xff0c;提供了 c11 foreach 接口 代码 vector<string> split(string s) {vector<string> res;const char *p strtok((char *) s.c_str(),…