Mqtt消费端实现的几种方式

2024-09-03 22:36
文章标签 实现 方式 几种 mqtt 消费

本文主要是介绍Mqtt消费端实现的几种方式,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

此处测试的mqtt的Broker是使用的EMQX 5.7.1,可移步至https://blog.csdn.net/tiantang_1986/article/details/140443513查看详细介绍

一、方式1

添加必要的依赖

<dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-integration</artifactId>
</dependency>
<dependency><groupId>org.springframework.integration</groupId><artifactId>spring-integration-mqtt</artifactId>
</dependency>

配置

# mqtt 服务端配置
spring:# mqtt 配置mqtt:url: tcp://127.0.0.1:1883,tcp://127.0.0.2:1883clientId: "00000001"       # 客户端Id(不可重复)username: <访问用户名>      # 认证的用户名password: <访问密码>        # 认证的密码qos: 1topic: test/#              # 监听的topic

读取配置文件

import org.apache.commons.lang3.StringUtils;
import lombok.Data;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.integration.mqtt.core.DefaultMqttPahoClientFactory;
import org.springframework.integration.mqtt.core.MqttPahoClientFactory;@Data
@Configuration
@ConfigurationProperties(prefix = "spring.mqtt")
public class MqttConfig {private String username;private String password;private String url;private String clientId;private String topic;private Integer qos;@Beanpublic MqttPahoClientFactory mqttClientFactory() {DefaultMqttPahoClientFactory factory = new DefaultMqttPahoClientFactory();MqttConnectOptions options = new MqttConnectOptions();options.setUserName(username);options.setPassword(password.toCharArray());if (StringUtils.isNotBlank(url) && url.contains(",")) {options.setServerURIs(url.split(","));} else {options.setServerURIs(new String[]{url});}        options.setCleanSession(true);//自动重连options.setAutomaticReconnect(true);//设置超时时间,单位为秒options.setConnectionTimeout(0);//设置心跳时间 单位为秒,表示服务器每隔 1.5*20秒的时间向客户端发送心跳判断客户端是否在线options.setKeepAliveInterval(90);//设置遗嘱消息options.setWill("will_topic", (this.clientId + "与服务器断开连接").getBytes(), qos, false);factory.setConnectionOptions(options);factory.setPersistence(new MemoryPersistence());return factory;}
}

MQTT消息入站配置

@Slf4j
@Configuration
@IntegrationComponentScan
public class MqttInboundConfiguration {@Resourceprivate MqttConfig mqttConfig;@Resourceprivate MqttPahoClientFactory mqttClientFactory;@Resourceprivate MqttMessageReceiver mqttMessageReceiver;@Beanpublic MessageChannel mqttInBoundChannel() {return new PublishSubscribeChannel();}@Beanpublic MessageProducerSupport mqttInbound() {MqttPahoMessageDrivenChannelAdapter adapter = new MqttPahoMessageDrivenChannelAdapter(mqttConfig.getClientId(), mqttClientFactory, mqttConfig.getTopic());DefaultPahoMessageConverter converter = new DefaultPahoMessageConverter();//传输Hex数据,如果是String则可使用默认值falseconverter.setPayloadAsBytes(true);adapter.setConverter(converter);adapter.setRecoveryInterval(10000);adapter.setQos(mqttConfig.getQos());adapter.setOutputChannel(mqttInBoundChannel());return adapter;}@Bean@ServiceActivator(inputChannel = "mqttInBoundChannel")public MessageHandler mqttMessageHandler() {return this.mqttMessageReceiver;}
}

消费者

@Slf4j
@Component
public class MqttMessageReceiver implements MessageHandler {@Resourceprivate DataConvertStrategyFactory convertStrategyContext;@Overridepublic void handleMessage(Message<?> message) throws MessagingException {MessageHeaders headers = message.getHeaders();String topic = (String) headers.get(MqttHeaders.RECEIVED_TOPIC);if (StringUtils.isNotBlank(topic)) {return;}byte[] payload = (byte[]) message.getPayload();log.info("topic: {}, message: {}", topic, HexUtils.bytesToHex(payload));//从topic中获取clientId,topic的格式:{业务}/{clientId}/{事件标识}Map<String, String> map = MqttDataConverter.covertTopic(topic);String clientId = map.get("clientId");log.info("clientId: {}", clientId);//topic中的事件标识String eventUrl = map.get("event");//自定义的enum,主要用来消息处理消息分组,相同组可以使用相同的数据转换服务Event[] events = Event.values();String deviceId = clientId;Arrays.stream(events).filter(item -> item.getEvent().equals(eventUrl)).findFirst().ifPresent(item -> {//使用策略模式实现DataConvertService convertService = convertStrategyContext.getStrategy(item.getGroup());convertService.convert(deviceId, eventUrl, payload);});}
}

数据转换服务接口,具体的数据解析只要实现这个接口就行

public interface DataConvertService {/*** 转换数据** @param clientId 设备SN* @param topic  topic* @param data  数据* @return*/Boolean convert(String clientId, String topic, byte[] data);/*** 获取转换器** @return*/String getConverter();
}

MQTT数据转换策略工厂

@Component
public class DataConvertStrategyFactory implements InitializingBean {@Resourceprivate List<DataConvertService> handlers;private Map<String, DataConvertService> dataConvertServiceMap = new ConcurrentHashMap<>();/*** 初始化*/@Overridepublic void afterPropertiesSet() {//进行初始化if (CollectionUtils.isNotEmpty(handlers)) {handlers.forEach(item -> {dataConvertServiceMap.put(item.getConverter(), item);});}}/*** 返回实际处理对象** @param strategy 处理策略* @return 实际处理对象*/public DataConvertService getStrategy(String strategy) {return dataConvertServiceMap.get(strategy);}
}

二、方式2

使用EMQXWebhook钩子
首先创建钩子函数,把需要监听的事件加上处理逻辑,示例:

@Slf4j
@RequestMapping("/mqtt/client")
@RestController
public class ClientController {@PostMapping("/webhook")public Result webhook(@RequestBody Map<String, Object> message) {log.info("webhook map:{}", message);String action = (String) message.get("action");String clientid = (String) message.get("clientid");if ("client_connected".equals(action)) {log.info("client:{} 上线", clientid);}if ("client_disconnected".equals(action)) {log.info("client:{} 下线", clientid);}if ("message.publish".equals(action)) {log.info("已接收到 client:{} 的消息:{}", clientid, message.get("payload"));}return Result.success("OK");}
}

然后在EMQX的Dashboard中创建Webhook,可以选择多个触发器
在这里插入图片描述
填好URL后可以进行测试,之后使用MQTTX进行消息发送测试
在这里插入图片描述
控制台输出日志
在这里插入图片描述

三、方式3

package com.iinplus.mqtt.handler;import com.iinplus.mqtt.config.MqttConfig;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.paho.client.mqttv3.*;
import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.stereotype.Component;import javax.annotation.Resource;@Slf4j
@Component
public class MqttSubscriber implements InitializingBean {@Resourceprivate MqttConfig config;@Overridepublic void afterPropertiesSet() {try {MqttClient client = new MqttClient(config.getUrl(), config.getClientId(), new MemoryPersistence());MqttConnectOptions options = new MqttConnectOptions();options.setUserName(config.getUsername());options.setPassword(config.getPassword().toCharArray());options.setCleanSession(true);options.setAutomaticReconnect(true);options.setConnectionTimeout(0);client.connect(options);client.subscribe(config.getTopic());//设置消息回调client.setCallback(new MqttMsgHandler());} catch (MqttException e) {log.error("MqttException:", e);}}
}

消息回调处理

@Slf4j
public class MqttMsgHandler implements MqttCallback {@Overridepublic void connectionLost(Throwable t) {// 连接丢失log.info("Connection lost:", t);}@Overridepublic void messageArrived(String topic, MqttMessage message) {// 接收到消息log.info("Message arrived:" + new String(message.getPayload()));}@Overridepublic void deliveryComplete(IMqttDeliveryToken token) {// 消息发送成功log.info("Delivery complete");}
}

这篇关于Mqtt消费端实现的几种方式的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/1134238

相关文章

hdu1043(八数码问题,广搜 + hash(实现状态压缩) )

利用康拓展开将一个排列映射成一个自然数,然后就变成了普通的广搜题。 #include<iostream>#include<algorithm>#include<string>#include<stack>#include<queue>#include<map>#include<stdio.h>#include<stdlib.h>#include<ctype.h>#inclu

【C++】_list常用方法解析及模拟实现

相信自己的力量,只要对自己始终保持信心,尽自己最大努力去完成任何事,就算事情最终结果是失败了,努力了也不留遗憾。💓💓💓 目录   ✨说在前面 🍋知识点一:什么是list? •🌰1.list的定义 •🌰2.list的基本特性 •🌰3.常用接口介绍 🍋知识点二:list常用接口 •🌰1.默认成员函数 🔥构造函数(⭐) 🔥析构函数 •🌰2.list对象

【Prometheus】PromQL向量匹配实现不同标签的向量数据进行运算

✨✨ 欢迎大家来到景天科技苑✨✨ 🎈🎈 养成好习惯,先赞后看哦~🎈🎈 🏆 作者简介:景天科技苑 🏆《头衔》:大厂架构师,华为云开发者社区专家博主,阿里云开发者社区专家博主,CSDN全栈领域优质创作者,掘金优秀博主,51CTO博客专家等。 🏆《博客》:Python全栈,前后端开发,小程序开发,人工智能,js逆向,App逆向,网络系统安全,数据分析,Django,fastapi

让树莓派智能语音助手实现定时提醒功能

最初的时候是想直接在rasa 的chatbot上实现,因为rasa本身是带有remindschedule模块的。不过经过一番折腾后,忽然发现,chatbot上实现的定时,语音助手不一定会有响应。因为,我目前语音助手的代码设置了长时间无应答会结束对话,这样一来,chatbot定时提醒的触发就不会被语音助手获悉。那怎么让语音助手也具有定时提醒功能呢? 我最后选择的方法是用threading.Time

Android实现任意版本设置默认的锁屏壁纸和桌面壁纸(两张壁纸可不一致)

客户有些需求需要设置默认壁纸和锁屏壁纸  在默认情况下 这两个壁纸是相同的  如果需要默认的锁屏壁纸和桌面壁纸不一样 需要额外修改 Android13实现 替换默认桌面壁纸: 将图片文件替换frameworks/base/core/res/res/drawable-nodpi/default_wallpaper.*  (注意不能是bmp格式) 替换默认锁屏壁纸: 将图片资源放入vendo

内核启动时减少log的方式

内核引导选项 内核引导选项大体上可以分为两类:一类与设备无关、另一类与设备有关。与设备有关的引导选项多如牛毛,需要你自己阅读内核中的相应驱动程序源码以获取其能够接受的引导选项。比如,如果你想知道可以向 AHA1542 SCSI 驱动程序传递哪些引导选项,那么就查看 drivers/scsi/aha1542.c 文件,一般在前面 100 行注释里就可以找到所接受的引导选项说明。大多数选项是通过"_

C#实战|大乐透选号器[6]:实现实时显示已选择的红蓝球数量

哈喽,你好啊,我是雷工。 关于大乐透选号器在前面已经记录了5篇笔记,这是第6篇; 接下来实现实时显示当前选中红球数量,蓝球数量; 以下为练习笔记。 01 效果演示 当选择和取消选择红球或蓝球时,在对应的位置显示实时已选择的红球、蓝球的数量; 02 标签名称 分别设置Label标签名称为:lblRedCount、lblBlueCount

Android平台播放RTSP流的几种方案探究(VLC VS ExoPlayer VS SmartPlayer)

技术背景 好多开发者需要遴选Android平台RTSP直播播放器的时候,不知道如何选的好,本文针对常用的方案,做个大概的说明: 1. 使用VLC for Android VLC Media Player(VLC多媒体播放器),最初命名为VideoLAN客户端,是VideoLAN品牌产品,是VideoLAN计划的多媒体播放器。它支持众多音频与视频解码器及文件格式,并支持DVD影音光盘,VCD影

webm怎么转换成mp4?这几种方法超多人在用!

webm怎么转换成mp4?WebM作为一种新兴的视频编码格式,近年来逐渐进入大众视野,其背后承载着诸多优势,但同时也伴随着不容忽视的局限性,首要挑战在于其兼容性边界,尽管WebM已广泛适应于众多网站与软件平台,但在特定应用环境或老旧设备上,其兼容难题依旧凸显,为用户体验带来不便,再者,WebM格式的非普适性也体现在编辑流程上,由于它并非行业内的通用标准,编辑过程中可能会遭遇格式不兼容的障碍,导致操

Kubernetes PodSecurityPolicy:PSP能实现的5种主要安全策略

Kubernetes PodSecurityPolicy:PSP能实现的5种主要安全策略 1. 特权模式限制2. 宿主机资源隔离3. 用户和组管理4. 权限提升控制5. SELinux配置 💖The Begin💖点点关注,收藏不迷路💖 Kubernetes的PodSecurityPolicy(PSP)是一个关键的安全特性,它在Pod创建之前实施安全策略,确保P