SpringBoot通过Moquette集成MQTT Broker快速实现设备管理消息订阅与发布

发布时间:2026/10/1 14:21:00

SpringBoot通过Moquette集成MQTT Broker快速实现设备管理消息订阅与发布
一、引言1.1 什么是 MQTTMQTT 是一种基于发布/订阅范式的轻量级消息传输协议。该协议专为计算资源受限、带宽有限以及网络连接不稳定的设备设计具有报文结构简洁、功耗低、传输可靠等特点。1.2 为什么选择 Moquette当业务需求比较简单没必要或开发受限无法在服务器部署第三方服务例如 EMQX、Mosquitto、HiveMQ那么就需要在 SpringBoot 中集成 MQTT Broker。Moquette 是一个轻量级嵌入式 Broker可以轻松封装在其他应用程序中。但由于0.18之前的版本不支持 MQTT5所以这里选择 Fraunhofer 维护的新分支。二、依赖配置!-- MQTT BrokerMoquette --dependencygroupIdde.fraunhofer.iosb.io.moquette/groupIdartifactIdmoquette-broker/artifactIdversion0.18.4/version/dependency三、代码整合3.1 创建配置类importio.moquette.broker.Server;importio.moquette.broker.config.IConfig;importio.moquette.broker.config.MemoryConfig;importio.moquette.interception.InterceptHandler;importjakarta.annotation.PostConstruct;importjakarta.annotation.PreDestroy;importorg.springframework.beans.factory.annotation.Autowired;importorg.springframework.context.annotation.Bean;importorg.springframework.context.annotation.Configuration;importjava.util.Arrays;importjava.util.List;importjava.util.Properties;ConfigurationpublicclassMqttBrokerConfig{privateServermqttBroker;AutowiredprivateCustomInterceptHandlerinterceptHandler;PostConstructpublicvoidstartBroker(){mqttBrokernewServer();PropertiespropsnewProperties();props.setProperty(port,1883);//监听端口默认 1883props.setProperty(host,0.0.0.0);//绑定地址0.0.0.0 表示允许所有 IP 访问props.setProperty(allow_anonymous,false);//禁止匿名登录必须用户名密码认证必须设为 false//性能相关配置props.setProperty(max_inflight_queue,1000);//单个客户端最大未确认消息队列长度防止消息堆积props.setProperty(netty.mqtt.message_size,65536);//最大消息体大小64KBprops.setProperty(persistence_enabled,true);//是否开启持久化props.setProperty(persistence_type,file);//持久化存储类型props.setProperty(data_path,D:\\haooo\\tempFile\\mqttFile);//持久化文件存放目录props.setProperty(persistence_recover_on_startup,true);//持久化启动时自动恢复props.setProperty(max_queued_messages,10000);//最大队列消息数IConfigconfignewMemoryConfig(props);try{ListInterceptHandlerhandlersArrays.asList(interceptHandler);mqttBroker.startServer(config,handlers,null,deviceAuthenticator,null);System.out.println(✅ Moquette MQTT Broker 启动成功端口 1883);}catch(Exceptione){System.err.println(❌ Moquette Broker 启动失败);e.printStackTrace();}}PreDestroypublicvoidstopBroker(){if(mqttBroker!null){mqttBroker.removeInterceptHandler(interceptHandler);mqttBroker.stopServer();System.out.println(✅ Moquette MQTT Broker 已停止);}}BeanpublicServermqttServer(){returnmqttBroker;//返回已启动的实例}}3.2 自定义认证importcom.github.benmanes.caffeine.cache.Cache;importcom.github.benmanes.caffeine.cache.Caffeine;importio.moquette.broker.security.IAuthenticator;importorg.springframework.stereotype.Component;importjava.time.Duration;/** * 自定义认证 * Moquette 支持插件式认证实现 IAuthenticator */ComponentpublicclassDeviceAuthenticatorimplementsIAuthenticator{privatefinalCacheString,BooleanauthCacheCaffeine.newBuilder().expireAfterWrite(Duration.ofMinutes(5)).build();OverridepublicbooleancheckValid(StringclientId,Stringusername,byte[]password){if(usernamenull||passwordnull||clientIdnull){System.out.println(❌ 设备认证失败用户名或密码或客户端标识为空);returnfalse;}StringpwdStrnewString(password);//验证设备是否存在5分钟内的验证结果进行缓存避免认证失败的设备频繁重试BooleancachedauthCache.getIfPresent(clientId:username);if(cached!null){returncached;}else{//数据库查询DevicedeviceDeviceMapper.selectDevice(username,clientId,pwdStr);if(device!null){System.out.println(✅ 设备认证成功: clientId | username | pwdStr);authCache.put(clientId:username,true);returntrue;}else{authCache.put(clientId:username,false);}}System.out.println(❌ 设备认证失败: clientId | username | pwdStr);returnfalse;}}3.3 消息处理importio.moquette.interception.AbstractInterceptHandler;importio.moquette.interception.messages.InterceptConnectMessage;importio.moquette.interception.messages.InterceptDisconnectMessage;importio.moquette.interception.messages.InterceptPublishMessage;importio.netty.buffer.ByteBuf;importorg.springframework.stereotype.Component;importjava.nio.charset.StandardCharsets;ComponentpublicclassCustomInterceptHandlerextendsAbstractInterceptHandler{OverridepublicStringgetID(){returnCustomInterceptHandler;}OverridepublicvoidonConnect(InterceptConnectMessagemsg){System.out.println(✅ 设备连接成功: msg.getClientID() | msg.getUsername());//TODO: 更新数据库设备在线离线状态}OverridepublicvoidonDisconnect(InterceptDisconnectMessagemsg){System.out.println(❌ 设备连接断开: msg.getClientID() | msg.getUsername());//TODO: 更新数据库设备在线离线状态}OverridepublicvoidonPublish(InterceptPublishMessagemsg){ByteBufbufmsg.getPayload();if(bufnull||buf.readableBytes()0){return;}byte[]bytesnewbyte[buf.readableBytes()];buf.getBytes(buf.readerIndex(),bytes);StringpayloadnewString(bytes,StandardCharsets.UTF_8);System.out.println(✅ 设备发布消息: msg.getClientID() | msg.getUsername() | msg.getTopicName() | payload);//处理设备发送的消息if(msg.getTopicName().matches(devices/[^/]/[^/]/resp)){//处理设备发布主题}else{System.out.println(❌ 未匹配设备发布消息: msg.getClientID() | msg.getUsername() | msg.getTopicName() | payload);}}OverridepublicvoidonSessionLoopError(Throwablecause){System.out.println(❌ 会话循环错误: cause.getMessage());}}3.4 发送指令importcom.alibaba.fastjson2.JSONObject;importio.moquette.broker.Server;importio.netty.buffer.ByteBuf;importio.netty.buffer.Unpooled;importio.netty.handler.codec.mqtt.*;importorg.springframework.beans.factory.annotation.Autowired;importorg.springframework.context.annotation.Lazy;importorg.springframework.stereotype.Service;ServicepublicclassMqttCommandService{AutowiredLazyprivateServermqttBroker;/** * 向指定设备发送指令 * param clientId 客户端标识 * param username 设备序列号 * param command 指令内容 */publicvoidsendCommand(StringclientId,Stringusername,JSONObjectcommand){Stringtopicdevices/clientId/username/cmd;//订阅主题try{MqttFixedHeadermqttFixedHeadernewMqttFixedHeader(MqttMessageType.PUBLISH,false,MqttQoS.AT_LEAST_ONCE,false,0);MqttPublishVariableHeadervariableHeadernewMqttPublishVariableHeader(topic,0);ByteBufpayloadUnpooled.wrappedBuffer(command.toJSONString().getBytes());MqttPublishMessagemsgnewMqttPublishMessage(mqttFixedHeader,variableHeader,payload);mqttBroker.internalPublish(msg,clientId);System.out.println(✅ 已下发指令到设备 username → command);}catch(Exceptione){e.printStackTrace();}}}四、测试使用 MQTTX 进行测试/** * 测试Controller */RestControllerRequestMapping(/test)publicclassTestControllerextendsBaseController{AutowiredprivateMqttCommandServicemqttCommandService;/** * MQTT消息发送 */PostMapping(/command)publicRcommand(RequestBodyJSONObjectparams){mqttCommandService.sendCommand(params.getString(clientId),params.getString(username),params.getJSONObject(command));returnsuccess();}}

相关新闻

Dify 流程节点配置完全指南

Dify 流程节点配置完全指南

2026/9/30 13:45:08

Dify 流程节点配置完全指南 文档说明: 本文档依据 Dify 官方文档整理,覆盖工作流、对话流全部节点配置,每个节点附带可直接复制使用的实操配置示例,适配 Version1.13.3 版本。 目录: 一、基础入口与输出节点 2 1.1 开始…

YOLOv11【第十九章:行业垂直应用与定制篇·第9节】电力巡检——输电线缺陷 + 绝缘子破损无人机检测!

YOLOv11【第十九章:行业垂直应用与定制篇·第9节】电力巡检——输电线缺陷 + 绝缘子破损无人机检测!

2026/9/26 19:54:04

🏆本文收录于专栏 《YOLOv11实战:从入门到深度优化》。 本专栏围绕 YOLOv11 的改进、训练、部署与工程优化 展开,系统梳理并复现当前主流的 YOLOv11 实战案例与优化方案,内容目前已覆盖 分类、检测、分割、追踪、关键点、OBB 检测 等多个方向。 整体坚持 持续更新 + 深度解…

zed 新特性:终端选中自动复制

zed 新特性:终端选中自动复制

2026/10/1 12:12:38

引言 对于后端、运维、云原生开发者而言,终端复制是日均高频操作:容器日志、镜像地址、SQL、文件路径、报错堆栈都需要反复选取粘贴。copy on select(选中自动复制)看似微小功能,却直接决定连续工作流是否被打断。Zed …

CANN/GE ACL数据集缓冲区添加函数

CANN/GE ACL数据集缓冲区添加函数

2026/9/29 22:00:59

aclmdlAddDatasetBuffer 【免费下载链接】ge GE(Graph Engine)是面向昇腾的图编译器和执行器,提供了计算图优化、多流并行、内存复用和模型下沉等技术手段,加速模型执行效率,减少模型内存占用。 GE 提供对 PyTorch、Te…

用ffmpeg高效批量调整图片尺寸的实战指南

用ffmpeg高效批量调整图片尺寸的实战指南

2026/9/30 20:35:35

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

Transformers 音频特征提取工具库 audio_utils 全解析:从 Mel 刻度换算到对数 Mel 频谱

Transformers 音频特征提取工具库 audio_utils 全解析:从 Mel 刻度换算到对数 Mel 频谱

2026/9/30 20:35:38

Transformers 音频特征提取工具库 audio_utils 全解析:从 Mel 刻度换算到对数 Mel 频谱 【免费下载链接】transformers 🤗 Transformers: the model-definition framework for state-of-the-art machine learning models in text, vision, audio, and mu…

RustFS 多节点集群重启与滚动升级实战:Readiness、Quorum 与 Degraded 模式完全指南

RustFS 多节点集群重启与滚动升级实战:Readiness、Quorum 与 Degraded 模式完全指南

2026/9/30 20:35:36

RustFS 多节点集群重启与滚动升级实战:Readiness、Quorum 与 Degraded 模式完全指南 【免费下载链接】rustfs 🚀2.3x faster than MinIO for 4KB object payloads. RustFS is an open-source, S3-compatible high-performance object storage system sup…

Java Integer缓存揭秘:128陷阱原理、避坑与面试全解

Java Integer缓存揭秘:128陷阱原理、避坑与面试全解

2026/9/30 20:35:33

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

RustFS Scanner 数据用量发布权威性决策:配额准入如何获得可用的权威依据

RustFS Scanner 数据用量发布权威性决策:配额准入如何获得可用的权威依据

2026/9/30 8:20:32

RustFS Scanner 数据用量发布权威性决策:配额准入如何获得可用的权威依据 【免费下载链接】rustfs 🚀2.3x faster than MinIO for 4KB object payloads. RustFS is an open-source, S3-compatible high-performance object storage system supporting mi…