一分钟在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.hqwc.cn/news/443068.html

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

相关文章

linux解决访问/下载github连接超时或下载慢的问题

问题 我这里是树莓派从github下载资源出现无法连接&#xff0c;连接超时的问题&#xff0c;如下所示解决方式 修改/etc/hosts文件 例&#xff1a; sudo nano /etc/hosts #添加如下 192.30.255.112 github.com git 185.31.16.184 github.global.ssl.fastly.net这里以树莓派为…

【C++】构造函数和析构函数详解

目录 前言 类中的六个默认成员函数 构造函数 概念 特性 析构函数 概念 特性&#xff1a; 前言 类中的六个默认成员函数 如果一个类中什么成员都没有&#xff0c;简称为空类。 空类中真的什么都没有吗&#xff1f;并不是&#xff0c;任何类在什么都不写时&#xff0c;编…

IT运维如何帮助企业降本增效?

IT监控运维管理技术发展应用和趋势 1、智能运维 随着人工智能和大数据技术的发展&#xff0c;智能运维将成为IT监控运维管理的重要趋势。通过利用机器学习、深度学习等技术&#xff0c;实现对IT系统的自动化监控、故障预测和智能维护&#xff0c;提高运维效率和质量。 2、容…

代码随想录——贪心算法

系列文章目录 代码随想录——回溯 代码随想录——贪心算法 文章目录 系列文章目录概述简单分发饼干摆动序列***最大子数组和买卖股票的最佳时机 II跳跃游戏跳跃游戏IIK次取反后最大化的数组和加油站***分发糖果柠檬水找零***根据身高重建队列全排列II 棋盘问题N皇后***解数独 …

C++ 数论相关题目 台阶-Nim游戏

现在&#xff0c;有一个 n 级台阶的楼梯&#xff0c;每级台阶上都有若干个石子&#xff0c;其中第 i 级台阶上有 ai 个石子(i≥1 )。 两位玩家轮流操作&#xff0c;每次操作可以从任意一级台阶上拿若干个石子放到下一级台阶中&#xff08;不能不拿&#xff09;。 已经拿到地面…

防御第五次作业-防火墙综合实验(nat、双机热备、带宽、选路)

目录 拓扑图 要求 1 2 3 4 5 办公区上网用户限制流量不超过60M 销售部限流 6 7 拓扑图 说明&#xff1a;基本配置完成。所有设备均可出网。 要求 1.办公区设备可以通过电信和移动两条链路上网&#xff0c;且需要保留一个公网ip不能用来转换。 2.分公司设备可以通…

flink cdc,standalone模式下,任务运行一段时间taskmanager挂掉

在使用flink cdc&#xff0c;配置任务运行&#xff0c;过了几天后&#xff0c;任务无故取消&#xff0c;超时&#xff0c;导致taskmanager挂掉&#xff0c;相关异常如下&#xff1a; 异常1&#xff1a; did not react to cancelling signal interrupting; it is stuck for 30 s…

快速排序|超详细讲解|入门深入学习排序算法

快速排序介绍 快速排序(Quick Sort)使用分治法策略。 它的基本思想是&#xff1a;选择一个基准数&#xff0c;通过一趟排序将要排序的数据分割成独立的两部分&#xff1b;其中一部分的所有数据都比另外一部分的所有数据都要小。然后&#xff0c;再按此方法对这两部分数据分别进…

第九节HarmonyOS 常用基础组件23-Menu、MenuItem、MenuItemGroup

1、描述 Menu&#xff1a;以垂直列表形式显示的菜单。 MenuItem&#xff1a;用来展示菜单Menu中具体的item菜单项。 MenuItemGroup&#xff1a;该组件用来展示菜单MenuItem的分组。 2、子组件 Menu&#xff1a;包含MenuItem、MenuItemGroup子组件。 MenuItem&#xff1a;…

第九节HarmonyOS 常用基础组件21-ImageAnimator

1、描述 提供帧动画组件来实现逐帧播放图片的能力&#xff0c;可以配置需要播放的图片列表&#xff0c;每张图片可以配置时长。 2、接口 ImageAnimator() 3、属性 参数名 参数类型 描述 images Array<ImageFrameInfo> 设置图片帧信息集合&#xff0c;每一帧的帧…

Oracle 的闪回技术是什么

什么是闪回 Oracle 数据库闪回技术是一组独特而丰富的数据恢复解决方案&#xff0c;能够有选择性地高效撤销一个错误的影响&#xff0c;从人为错误中恢复。闪回是一种数据恢复技术&#xff0c;它使得数据库可以回到过去的某个状态&#xff0c;可以满足用户的逻辑错误的快速恢复…

【数据结构与算法】之哈希表系列-20240130

这里写目录标题 一、383. 赎金信二、387. 字符串中的第一个唯一字符三、389. 找不同四、409. 最长回文串五、448. 找到所有数组中消失的数字六、594. 最长和谐子序列 一、383. 赎金信 简单 给你两个字符串&#xff1a;ransomNote 和 magazine &#xff0c;判断 ransomNote 能不…