springBoot集成emqx 实现mqtt消息的发送订阅

介绍

我们可以想象这么一个场景,我们java应用想要采集到电表a的每小时的用电信息,我们怎么拿到电表的数据?一般我们会想 直接 java 后台发送请求给电表,然后让电表返回数据就可以了,事实上,我们java应用发送请求请求电表的数据信息并不是发到电表上,而是发送到 服务端 (broker)上,请求服务器 给我们电表的信息,而电表会把数据 按照mqtt协议 源源不断的发送到服务端,服务端可以把数据存储到物联网数据库上,也可以由我们java应用手动存储到物联网数据库上

而我们怎么知道电表发送到服务端的哪里,java应用又怎么请求到该电表发送的位置?

这就 引出了 一个概念,主题 (topic) ,这个topic在mqtt中不需要手动的创建,只要又客户端订阅或者发布消息,主题就会被自动创建出来

而我们服务端用的最多的就是 集成好的emqx服务器,本文我们也用的是集成好的emqx的服务端

,我们先是 一个电表 订阅好一个固定的主题,然后 源源不断的往服务端发消息,然后我们java应用订阅这个主题,这样 java应用就能持续的拿到电表的数据了

具体什么是主题,主题怎么设置的 ,mqtt协议的具体协议内容,直接登录emqx官网查看即可

MQTT 最全教程:从入门到精通 | EMQ

而emqx服务器是怎么在linux系统上搭建的呢,具体直接看文档即可,输入文档对应的yum命令就可以直接 在linux服务器上安装了

在 CentOS/RHEL 上安装 EMQX | EMQX文档

文本主要书写代码的实现

代码实现

我们的yml文件如下

 我们后续 java应用订阅消息 都要到服务端 emqx 的1883端口

实体类如下

 

这里解释一下 clientid 是不固定的,随机的每一个发布/订阅消息的客户端都有一个唯一的clientid

而username 和password 是 客户端连接到 服务端的认证账户,多个客户端可以使用一个 账号密码

客户端代码实现

@Slf4j
@Component
@RequiredArgsConstructor
public class EMQXClient {
private  final MqttDefaultProperties mqttDefaultProperties;
private  final IMessageCallbackImpl mqttCallback;private IMqttClient mqttClient;/*** 初始化客户端对象*/
public  boolean initMqttClient(String clientId,String serverUrl)  {MemoryPersistence memoryPersistence = new MemoryPersistence();try {if(Objects.isNull(clientId)){clientId= mqttDefaultProperties.getDefaultClientId();}if(Objects.isNull(serverUrl)){serverUrl= mqttDefaultProperties.getServerUrl();}mqttClient = new MqttClient(serverUrl, clientId, memoryPersistence);} catch (MqttException e) {log.info("mqtt创建异常:{}", e.getMessage());return  false;}return true;
}public  boolean initMqttClient()  {MemoryPersistence memoryPersistence = new MemoryPersistence();try {mqttClient = new MqttClient(mqttDefaultProperties.getServerUrl(),mqttDefaultProperties.getDefaultClientId(),  memoryPersistence);} catch (MqttException e) {log.info("mqtt创建异常:{}", e.getMessage());return  false;}return true;}/*** 获取连接* @return*/public   boolean connect()  {MqttConnectOptions mqttConnectOptions = new MqttConnectOptions();//当客户端会话关闭的时候  对应的broker也关闭mqttConnectOptions.setCleanSession(true);//自动重连mqttConnectOptions.setAutomaticReconnect(true);mqttConnectOptions.setUserName(mqttDefaultProperties.getDefaultUserName());mqttConnectOptions.setPassword(mqttDefaultProperties.getDefaultPassword().toCharArray());mqttClient.setCallback(mqttCallback);try {mqttClient.connect(mqttConnectOptions);} catch (MqttException e) {log.info("客户端连接异常:{}",e.getMessage());return  false;}return true;}public   boolean connect(String username,String password)  {if(Objects.isNull(username)){username=mqttDefaultProperties.getDefaultUserName();}if(Objects.isNull(password)){password= mqttDefaultProperties.getDefaultPassword();}MqttConnectOptions mqttConnectOptions = new MqttConnectOptions();//当客户端会话关闭的时候  对应的broker也关闭mqttConnectOptions.setCleanSession(true);//自动重连mqttConnectOptions.setAutomaticReconnect(true);mqttConnectOptions.setUserName(username);mqttConnectOptions.setPassword(password.toCharArray());mqttClient.setCallback(mqttCallback);try {mqttClient.connect(mqttConnectOptions);} catch (MqttException e) {log.info("客户端连接异常:{}",e.getMessage());return  false;}return true;
}/*** 断开连接* @return*/public  boolean disConnect(){try {mqttClient.disconnect();} catch (MqttException e) {log.info("客户端断开连接异常:{}",e.getMessage());return  false;}
return true;
}/**** @param topic 主题* @param msg 消息内容* @param qosEnum* @param retain  新的订阅者来了是否能拿到之前的 最新的一次消息* @return*/public  boolean publish(String topic, String msg, QosEnum qosEnum,boolean retain){int uniqueInt = (int) (System.nanoTime() & 0xFFFFFFFFL);//取纳秒时间戳低32位MqttMessage mqttMessage = new MqttMessage();mqttMessage.setPayload(msg.getBytes());mqttMessage.setQos(qosEnum.getType());mqttMessage.setRetained(retain);mqttMessage.setId(uniqueInt);try {mqttClient.publish(topic,mqttMessage);} catch (MqttException e) {log.info("客户端发送消息失败:{}",e.getMessage());return  false;}return  true;}/**** @param topicFilters 要订阅的主题 例子 testtopic/#* @param qosEnum* @return*/public boolean subscribe(String topicFilters,QosEnum qosEnum){try {mqttClient.subscribe(topicFilters,qosEnum.getType());} catch (MqttException e) {log.info("订阅主题失败:{}",e.getMessage());return  false;}return true;}
public boolean unSubscribe(String topicFilter){try {mqttClient.unsubscribe( topicFilter);} catch (MqttException e) {log.info("取消订阅主题失败:{}",e.getMessage());return false;}return  true;
}}

我们着重关注的是

我们想连接 服务端 是不是得有 一个client ,那这个client就对应IMqttclient

,我们java应用客户端连接上服务端之后,是不是得订阅主题,订阅之后的逻辑在哪里,就在

IMessageCallbackImpl

这里面就是 书写的 客户端收到服务端发来的消息之后的处理情况

@Slf4j
@Component
public class IMessageCallbackImpl  implements MessageCallback {@Overridepublic void connectionLost(Throwable cause) {//丢失对服务端的连接后触发该方法回调,此处可以做一些特殊处理,比如重连 或者记录 日志之类的log.info("丢失了对broker的连接");}/*** 订阅到消息后的回调* 该方法由mqtt客户端同步调用,在此方法未正确返回之前,不会发送ack确认消息到broker* 一旦该方法向外抛出了异常客户端将异常关闭,当再次连接时;所有QoS1,QoS2且客户端未进行ack确认的消息都将由broker服务器再次发送到客户端* @param topic* @param message* @throws Exception*/@Overridepublic void messageArrived(String topic, MqttMessage message) throws Exception {log.info("订阅到了消息;topic={},messageid={},qos={},msg={}",topic,message.getId(),message.getQos(),new String(message.getPayload()));}/*** 消息发布完成且收到ack确认后的回调* QoS0:消息被服务端发出后触发一次* QoS1:当收到broker的PUBACK消息后触发* QoS2:当收到broer的PUBCOMP消息后触发* @param token*/@Overridepublic void deliveryComplete(IMqttDeliveryToken token) {int messageId = token.getMessageId();String[] topics = token.getTopics();log.info("消息发送完成,messageId={},topics={}",messageId,topics);}
}

 

我们用一个 bean 在初始化的时候就订阅一个主题,这样 只要有 客户端往主题上发消息,我们就能收到了

而我们这个时候 没有硬件,怎么办呢,很简单,直接下载一个mqttx 模拟硬件发送消息到主题,启动springboot,就能看到消息的发送与接收了

当然 这实现的紧紧是最简单的协议的发送接收,后面还有许多的高级功能等我们使用,具体的可以查阅官方文档

相关新闻

【含文档+PPT+源码】基于SpringBoot电脑DIY装机教程网站的设计与实现

【含文档+PPT+源码】基于SpringBoot电脑DIY装机教程网站的设计与实现

项目介绍 本课程演示的是一款 基于SpringBoot电脑DIY装机教程网站的设计与实现,主要针对计算机相关专业的正在做毕设的学生与需要项目实战练习的 Java 学习者。 1.包含:项目源码、项目文档、数据库脚本、软件工具等所有资料 2.带你从零开始部署运行本…

2026/7/18 3:04:19 阅读更多 →
用OpenCV写个视频播放器可还行?(Python版)

用OpenCV写个视频播放器可还行?(Python版)

引言 提到OpenCV,大家首先想到的可能是图像处理、目标检测,但你是否想过——用OpenCV实现一个带进度条、倍速播放、暂停功能的视频播放器?本文将通过一个实战项目,带你深入掌握OpenCV的视频处理能力,并解锁以下功能&a…

2026/7/19 2:40:08 阅读更多 →
让 LabVIEW 程序更稳定

让 LabVIEW 程序更稳定

LabVIEW 开发的系统,尤其是工业级应用,往往需要长时间稳定运行,容不得崩溃、卡顿或数据丢失。然而,许多系统在实际运行中会遭遇内存泄漏、通信中断、界面卡顿等问题,导致生产中断甚至设备损坏。如何设计一个既稳定又易…

2026/7/19 2:58:26 阅读更多 →
ZBrush2026.2.1免安装版全面解析:部署指南与性能优化

ZBrush2026.2.1免安装版全面解析:部署指南与性能优化

如果你是一名3D建模师或数字雕塑爱好者,一定对ZBrush这个名字不陌生。作为业界顶级的数字雕刻软件,ZBrush以其强大的笔刷系统和直观的雕刻体验,成为游戏、影视、动画行业不可或缺的工具。但传统安装过程的复杂性——激活失败、版本冲突、系统…

2026/7/19 20:59:46 阅读更多 →
芝柏官方服务项目及价格查询|网点地址及24小时电话权威信息通告(2026年7月最新) - 亨得利官方服务中心

芝柏官方服务项目及价格查询|网点地址及24小时电话权威信息通告(2026年7月最新) - 亨得利官方服务中心

芝柏官方售后服务严格遵循品牌总部制定的全国统一标准,为每一位腕表拥有者提供专业、透明、可追溯的保养与维修服务。作为拥有超过两百年制表历史的瑞士高级钟表品牌,芝柏对售后服务的技术规范、配件供应及服务流程均…

2026/7/19 20:59:46 阅读更多 →
Spring AI构建智能航空客服:RAG技术实现与Java工程实践

Spring AI构建智能航空客服:RAG技术实现与Java工程实践

这次我们来看一个基于 Spring AI 的智能航空客服项目。这个项目结合 Java 技术栈和大模型能力,通过 RAG 技术构建企业级航空客服系统,能够处理航班查询、退改签政策、行李托运等常见航空业务问题。项目最值得关注的是它完整的 RAG 流程:从文档…

2026/7/19 20:59:46 阅读更多 →
AM62L AES引擎寄存器配置实战:从硬件加速原理到嵌入式安全开发

AM62L AES引擎寄存器配置实战:从硬件加速原理到嵌入式安全开发

1. AM62L AES引擎:从寄存器手册到实战配置的深度解析如果你正在基于TI的AM62L处理器开发需要数据加密功能的产品,比如物联网网关、工业控制器或者支付终端,那么你大概率绕不开它的硬件AES引擎。手册里那几十页密密麻麻的寄存器描述&#xff0…

2026/7/19 20:59:46 阅读更多 →
2019年Android高级工程师面试核心考点与实战解析

2019年Android高级工程师面试核心考点与实战解析

1. 2019年Android高级工程师面试全景分析2019年的移动互联网行业正处于从增量市场向存量市场转型的关键节点,各大厂对Android工程师的要求发生了显著变化。当时我在准备美团和字节跳动的面试时,明显感受到技术考察的深度和广度都在提升。不同于初级工程师…

2026/7/19 21:00:47 阅读更多 →
AM62L MMR寄存器实战:ADC、EPWM与EQEP模块配置详解

AM62L MMR寄存器实战:ADC、EPWM与EQEP模块配置详解

1. 项目概述与MMR寄存器核心价值 在嵌入式开发,尤其是工业控制、电机驱动这类对实时性要求极高的领域,我们工程师与硬件打交道最直接、最底层的接口,就是内存映射寄存器。你可能听过很多次MMR,但真正动手配置时,面对动…

2026/7/19 21:00:47 阅读更多 →
鸿蒙 ArkTS 实战:Emoji Idiom Guess 从表情成语猜谜到交互闭环完整解析

鸿蒙 ArkTS 实战:Emoji Idiom Guess 从表情成语猜谜到交互闭环完整解析

鸿蒙 ArkTS 实战:Emoji Idiom Guess 从表情成语猜谜到交互闭环完整解析 前言 Emoji Idiom Guess 是一个基于鸿蒙 ArkTS 编写的单页互动应用,核心围绕 表情线索、答案输入、首字母提示和收藏关卡 展开。项目没有依赖复杂服务端,也没有把逻辑…

2026/7/19 0:00:03 阅读更多 →
Unity与Python本地通信:基于Flask的跨语言数据交换实战

Unity与Python本地通信:基于Flask的跨语言数据交换实战

1. 项目概述:为什么我们需要一个本地通信服务器?在游戏开发、数字孪生、仿真训练等众多领域,Unity作为强大的实时3D内容创作平台,其核心逻辑通常由C#驱动。然而,当我们需要进行复杂的数据分析、机器学习推理、科学计算…

2026/7/19 0:00:04 阅读更多 →
科研课题设计全流程:从选题到成果落地的实战指南

科研课题设计全流程:从选题到成果落地的实战指南

1. 课题设计全流程解析:从选题到成果落地的实战指南课题设计是科研工作者、高校师生以及企业研发人员日常工作中的核心环节。一个优秀的课题设计不仅决定了研究的方向和质量,更直接影响最终成果的学术价值和应用前景。作为在科研一线摸爬滚打多年的从业者…

2026/7/19 0:00:04 阅读更多 →
鸿蒙 ArkTS 实战:Emoji Idiom Guess 从表情成语猜谜到交互闭环完整解析

鸿蒙 ArkTS 实战:Emoji Idiom Guess 从表情成语猜谜到交互闭环完整解析

鸿蒙 ArkTS 实战:Emoji Idiom Guess 从表情成语猜谜到交互闭环完整解析 前言 Emoji Idiom Guess 是一个基于鸿蒙 ArkTS 编写的单页互动应用,核心围绕 表情线索、答案输入、首字母提示和收藏关卡 展开。项目没有依赖复杂服务端,也没有把逻辑…

2026/7/19 0:00:03 阅读更多 →
Unity与Python本地通信:基于Flask的跨语言数据交换实战

Unity与Python本地通信:基于Flask的跨语言数据交换实战

1. 项目概述:为什么我们需要一个本地通信服务器?在游戏开发、数字孪生、仿真训练等众多领域,Unity作为强大的实时3D内容创作平台,其核心逻辑通常由C#驱动。然而,当我们需要进行复杂的数据分析、机器学习推理、科学计算…

2026/7/19 0:00:04 阅读更多 →
科研课题设计全流程:从选题到成果落地的实战指南

科研课题设计全流程:从选题到成果落地的实战指南

1. 课题设计全流程解析:从选题到成果落地的实战指南课题设计是科研工作者、高校师生以及企业研发人员日常工作中的核心环节。一个优秀的课题设计不仅决定了研究的方向和质量,更直接影响最终成果的学术价值和应用前景。作为在科研一线摸爬滚打多年的从业者…

2026/7/19 0:00:04 阅读更多 →
ai agent框架spring ai/alibaba 源码原理分析(六) agent和组件

ai agent框架spring ai/alibaba 源码原理分析(六) agent和组件

简介 saa是java的ai agent框架,本系列将深入剖析 Spring AI Alibaba 的源码实现与核心原理,不仅可以指导agent的开发,更可以改造框架,增加新特性 系列内容: 系列(一) 架构 完成 系列(三) 调用 I 工具 完成 II M…

2026/7/19 0:01:20 阅读更多 →
终极指南:如何用Steam-auto-crack实现Steam游戏自动破解

终极指南:如何用Steam-auto-crack实现Steam游戏自动破解

终极指南:如何用Steam-auto-crack实现Steam游戏自动破解 【免费下载链接】Steam-auto-crack Steam Game Automatic Cracker 项目地址: https://gitcode.com/gh_mirrors/st/Steam-auto-crack Steam-auto-crack是一款功能强大的Steam游戏自动破解工具&#xff…

2026/7/19 9:10:31 阅读更多 →
移动端游戏功耗测试实战:电流、功率、亮度和场景对比

移动端游戏功耗测试实战:电流、功率、亮度和场景对比

移动端游戏功耗测试:先控制变量,再比较优化是否真的省电 摘要:功耗测试最容易犯的错误,是拿两次不同温度、不同亮度、不同场景的平均功率直接比较。本文给出一套可复现的游戏功耗测试方法,覆盖引擎特性验证、版本回归和黑盒体验测试,并说明如何把功耗与帧率、温控、CPU/G…

2026/7/19 19:29:48 阅读更多 →