资讯详情

资讯详情

基于 Netty 与 MQTT 的 Spring Boot 物联网消息服务端与客户端实现

简介一份基于Netty、MQTT 3.1.1协议、Spring Boot与JDK 8实现的MQTT服务器与客户端项目面向物联网通信学习者和Java后端开发者尤其适合作为毕业设计参考。资源共113个文件以87个Java源码为核心辅以XML配置、YAML配置、JKS证书、factories文件及说明文档压缩包仅153KB目录结构清晰便于整体研读。已有536人学习下载。通过该资源可深入理解MQTT连接、发布订阅、会话保持机制掌握Netty异步网络编程与Spring Boot集成方式代码涵盖服务端、客户端及测试模块并附授权码说明与部署指引。可对照源码和教程学习实际开发中网络通信、状态管理、异常处理等问题帮助快速搭建可运行的MQTT通信Demo提升Java后端与IoT实战能力。1. 基于 Netty MQTT 3.1.1 Spring Boot JDK8 实现 MQTT 服务端跟客户端到底解决什么问题基于 Netty MQTT 3.1.1 Spring Boot JDK8 写一整套服务端和客户端听上去像重复造轮子但真正把netty-codec-mqtt的MqttDecoder接进 Spring Boot 时你会发现难的不是解析 CONNECT/PUBLISH而是服务端要自己维护连接会话、订阅关系和 retained message。这个组合的典型应用是物联网充电桩接入设备端内存只有几十 KB跑不了重量级协议MQTT 3.1.1 刚好够用后端用 Spring Boot 进程同时承担 REST API 和 MQTT 网关JDK8 部署环境也最容易找。适合想弄懂 MQTT 服务器搭建过程、又不满足于直接改 mosquitto 配置的开发者也适合需要把 Broker 能力嵌进业务系统而不是单独部署一套中间件的团队。2. 把 MQTT 3.1.1 控制报文和 Netty pipeline 对齐后面代码才有底一套 MQTT 服务端和客户端本质上就是在两条 Netty channel 里跑同一套报文处理逻辑。先看协议层再写 handler否则你会被MqttDecoder的返回类型绕晕。2.1 固定头、剩余长度和报文类型连接帧 2 分钟看明白MQTT 3.1.1 的报文分三部分固定头、可变头、载荷。固定头第一字节的高四位是报文类型低四位是 DUP、QoS、RETAIN 标志位接着是变长编码的剩余长度最多 4 个字节。具体的报文类型可以对照这张表。报文类型值方向用途CONNECT1C → S客户端发起连接CONNACK2S → C服务端确认连接结果PUBLISH3双向发布应用消息PUBACK4双向QoS 1 确认PUBREC / PUBREL / PUBCOMP5 / 6 / 7双向QoS 2 四次握手SUBSCRIBE8C → S订阅主题过滤器SUBACK9S → C订阅确认PINGREQ / PINGRESP12 / 13双向心跳DISCONNECT14C → S断开连接Netty 的netty-codec-mqtt模块已经把这些报文封装成了MqttMessage的子类型不需要自己拼字节。但你要知道解码边界在哪。剩余长度使用 128 进制变长编码MqttDecoder正是靠它判断一个完整报文结束的位置。int first buf.readUnsignedByte(); int messageType (first 4) 0x0F; int flags first 0x0F; int remainingLength 0; int multiplier 1; int digit 0; do { digit buf.readUnsignedByte(); remainingLength (digit 127) * multiplier; multiplier * 128; } while ((digit 128) ! 0);这段代码从左到右读出固定头第一字节再连续读剩余长度字节直到最高位为 0。MqttDecoder内部就是这样计算报文边界遇到半包会缓存并等待下一个 TCP segment 到达。所以你在 Netty pipeline 里只要把它放在最前面不需要自己处理粘包拆包。2.2 把 MQTT 解码放进 Netty 的 ChannelPipeline协议级粘包拆包在 Spring Boot 里启动 Broker 时最外层是 TCP 连接然后是 MQTT 编解码器、空闲检测、业务 handler。常见做法是ChannelPipeline pipeline ch.pipeline(); pipeline.addLast(decoder, new MqttDecoder(1024 * 1024)); pipeline.addLast(encoder, MqttEncoder.INSTANCE); pipeline.addLast(idle, new IdleStateHandler(60, 0, 0)); pipeline.addLast(broker, new MqttBrokerHandler(...));MqttDecoder的构造参数是maxBytesInMessage单个 MQTT 报文超过 1MB 会直接抛异常防止恶意连接消耗内存。MqttEncoder.INSTANCE是单例编码器把MqttMessage写回 ByteBuf。IdleStateHandler的 60 秒只作用于读空闲如果 60 秒任何报文都没到会触发IdleStateEvent再由后面的 handler 关闭连接。心跳判断的逻辑放在userEventTriggered里这是 Netty 处理用户自定义事件的标准入口。选型理由很简单Netty 的 EventLoop 模型让每个 Channel 绑定一个线程MQTT 长连接数量上去了也不会阻塞 HTTP 线程JDK8 的ConcurrentHashMap、Lambda、CompletableFuture都能很顺手地用在会话表、消息转发这些场景。这里不选 Mosquitto 是因为业务方需要把设备在线状态、订阅关系直接写进 Spring Boot 的缓存和数据库内嵌一套 Netty Broker 可以把这部分逻辑合并成一个进程。3. 用 Spring Boot 装配 Netty 实现 MQTT 服务端依赖、初始化和会话路由服务端要承接三类设备只上报状态的电表、下发指令的充电桩、订阅告警的运维后台。下面的实现方式对这三类都够用重点在 Spring 容器和 Netty 生命周期不能打架。3.1 Maven 依赖和 application.ymlJDK8 环境里 Spring Boot 版本怎么锁JDK8 对应 Spring Boot 2.x不要强行用 Spring Boot 3.x否则最低要求 JDK17 启动就报错。一般锁在 2.7.18 这个最终版本依赖管理已经足够稳定。就算以后升级到 Spring Boot 3.x Netty MQTT 做物联网智能充电桩这部分代码结构也基本不变。parent groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-parent/artifactId version2.7.18/version /parent properties java.version1.8/java.version /properties dependencies dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter/artifactId /dependency dependency groupIdio.netty/groupId artifactIdnetty-codec-mqtt/artifactId /dependency dependency groupIdio.netty/groupId artifactIdnetty-transport/artifactId /dependency /dependenciesspring-boot-starter不带 Web 容器如果你同时要提供 HTTP API再单独加spring-boot-starter-web但内嵌 Tomcat 和 Netty 各管各的端口。配置集中放在 application.ymlmqtt: broker: port: 1883 boss-threads: 1 worker-threads: 0 heartbeat-timeout: 60worker-threads 为 0 表示使用 Netty 默认值按 CPU 核心数乘 2。这个参数在大规模连接场景下需要调大但不建议一开始就写死。3.2 Netty 服务端启动与 Spring 生命周期绑定stop 要比容器慢一步如果用PostConstruct启动 Netty服务销毁时EventLoopGroup可能还没释放Spring 容器就把依赖它的 Bean 关掉了出现端口未释放或线程泄漏。我一般会实现SmartLifecycle让 Netty 的启动阶段排在一般 Bean 之后关闭阶段排在一般 Bean 之前。Component public class MqttBrokerServer implements SmartLifecycle { private final MqttBrokerHandler brokerHandler; private EventLoopGroup bossGroup; private EventLoopGroup workerGroup; private Channel serverChannel; private volatile boolean running; Override public void start() { bossGroup new NioEventLoopGroup(1); workerGroup new NioEventLoopGroup(0); ServerBootstrap b new ServerBootstrap(); b.group(bossGroup, workerGroup) .channel(NioServerSocketChannel.class) .childHandler(new ChannelInitializerSocketChannel() { Override protected void initChannel(SocketChannel ch) { ChannelPipeline p ch.pipeline(); p.addLast(decoder, new MqttDecoder(1024 * 1024)); p.addLast(encoder, MqttEncoder.INSTANCE); p.addLast(idle, new IdleStateHandler(60, 0, 0)); p.addLast(broker, brokerHandler); } }) .option(ChannelOption.SO_BACKLOG, 1024) .childOption(ChannelOption.TCP_NODELAY, true); try { serverChannel b.bind(port).sync().channel(); running true; } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new RuntimeException(MQTT broker bind failed, e); } } Override public void stop() { if (serverChannel ! null) { serverChannel.close(); } bossGroup.shutdownGracefully(); workerGroup.shutdownGracefully(); running false; } Override public boolean isRunning() { return running; } }绑定端口成功后才置running true这样 Spring Boot 的健康检查不会误判服务可用。option(ChannelOption.SO_BACKLOG, 1024)是 accept 队列长度设备集中上报时可以再调大TCP_NODELAY关闭 Nagle 算法让下发指令不因为等待小包合并而延迟。3.3 Handler 里维护会话CONNACK、SUBSCRIBE、PUBLISH 路由不丢消息MqttBrokerHandler必须加ChannelHandler.Sharable否则每个新连接都会创建一个 handler 实例全局会话表就存不住了。下面用传统 switch 写法因为 JDK8 不支持箭头语法。ChannelHandler.Sharable Component public class MqttBrokerHandler extends ChannelInboundHandlerAdapter { private final MapString, Channel clients new ConcurrentHashMap(); private final MapString, SetString subscriptions new ConcurrentHashMap(); private final MapString, RetainedMessage retained new ConcurrentHashMap(); Override public void channelRead(ChannelHandlerContext ctx, Object msg) { MqttMessage message (MqttMessage) msg; MqttMessageType messageType message.fixedHeader().messageType(); switch (messageType) { case CONNECT: handleConnect(ctx, (MqttConnectMessage) message); break; case SUBSCRIBE: handleSubscribe(ctx, (MqttSubscribeMessage) message); break; case PUBLISH: handlePublish(ctx, (MqttPublishMessage) message); break; case PINGREQ: ctx.writeAndFlush(new MqttMessage(new MqttFixedHeader( MqttMessageType.PINGRESP, false, MqttQoS.AT_MOST_ONCE, false, 0))); break; case DISCONNECT: ctx.close(); break; default: break; } } Override public void userEventTriggered(ChannelHandlerContext ctx, Object evt) throws Exception { if (evt instanceof IdleStateEvent) { ctx.close(); } else { super.userEventTriggered(ctx, evt); } } }handleConnect里拿msg.payload().clientIdentifier()做 key存Channel回 CONNACKhandleSubscribe把 topic filter 加入订阅表回 SUBACKhandlePublish根据主题匹配订阅者QoS 1 还要回 PUBACK。整体动作可以压缩成一张处理表。入站报文服务端动作出站报文CONNECT校验 clientId 和用户名密码CONNACKSUBSCRIBE按 filter 记录订阅关系SUBACKPUBLISH匹配订阅者并转发PUBLISH / PUBACKPINGREQ重置连接空闲计时PINGRESPDISCONNECT清理会话。无每次转发都直接拿Channel写不走消息队列。如果消息量和订阅关系复杂再在handlePublish和订阅表之间加一个内存队列或者 Redis topic这个结构不用大改。4. 用 Netty 写一个 MQTT 客户端来验证服务端QoS 和心跳按报文级检查服务端写完了光靠教科书上的状态机不能说服自己。用 Netty 写一个最小客户端直接看 CONNACK 和 SUBACK 的返回码比任何模拟器都直观。4.1 最小客户端连接参数、CONNECT 包和心跳定时客户端 handler 的思路和服务端对称Bootstrap连接pipeline 里同样是MqttDecoder、MqttEncoder、IdleStateHandler。这里用写空闲 10 秒触发心跳跟服务端的 60 秒读空闲配合。Bootstrap b new Bootstrap(); b.group(new NioEventLoopGroup(1)) .channel(NioSocketChannel.class) .handler(new ChannelInitializerSocketChannel() { Override protected void initChannel(SocketChannel ch) { ChannelPipeline p ch.pipeline(); p.addLast(decoder, new MqttDecoder(1024 * 1024)); p.addLast(encoder, MqttEncoder.INSTANCE); p.addLast(idle, new IdleStateHandler(0, 10, 0)); p.addLast(handler, new MqttClientHandler()); } }); Channel ch b.connect(127.0.0.1, 1883).sync().channel(); MqttFixedHeader connectHeader new MqttFixedHeader( MqttMessageType.CONNECT, false, MqttQoS.AT_MOST_ONCE, false, 0); MqttConnectVariableHeader variable new MqttConnectVariableHeader( MQTT, 4, false, false, false, MqttQoS.AT_MOST_ONCE, false, true, 60); MqttConnectPayload payload new MqttConnectPayload(demo-client, null, null, null, null); ch.writeAndFlush(new MqttConnectMessage(connectHeader, variable, payload));MqttConnectVariableHeader的参数依次是协议名MQTT、协议版本 4、是否有用户名、是否有密码、是否有遗嘱标志、遗嘱 QoS、是否保留遗嘱、是否 cleanSession、心跳秒数。这里 cleanSession 为 true服务端不会在断线后保留会话重连后订阅关系需要重建适合客户端刚联调的场景。心跳发送写在MqttClientHandler.userEventTriggered里Override public void userEventTriggered(ChannelHandlerContext ctx, Object evt) { if (evt instanceof IdleStateEvent) { MqttFixedHeader ping new MqttFixedHeader( MqttMessageType.PINGREQ, false, MqttQoS.AT_MOST_ONCE, false, 0); ctx.writeAndFlush(new MqttMessage(ping)); } else { super.userEventTriggered(ctx, evt); } }4.2 用客户端和外部工具联调 QoS 与保留消息服务端的 QoS 降级也要在联调时验证发送方用 QoS 1接收方订阅时申请 QoS 0最终投递 QoS 取两者最小值。这个规则不在 Netty 编解码层而在 handler 的转发逻辑里弄错会出现订阅端收不到 PUBACK 或重复消息。QoS 级别语义适用场景0最多一次不确认温湿度传感器高频上报1至少一次PUBACK 确认充电桩状态切换允许重复2恰好一次四次握手计费指令重复会导致乱账联调时除了自己写的 Netty 客户端我还会用标准 MQTT 客户端软件和命令行工具交叉验证。常见的 mqtt 客户端软件都能连本机 1883 端口服观察报文命令行更快mosquitto_sub -h 127.0.0.1 -p 1883 -t evse/device/001/# -q 1 mosquitto_pub -h 127.0.0.1 -p 1883 -t evse/device/001/status -m charging -q 1如果mosquitto_sub能收到charging而 Netty 客户端收不到问题在订阅表匹配如果两边都收不到先抓包看 CONNACK 是不是返回了 1 或者 2而不是直接改转发逻辑。mqtt 怎么连接这个问题十次里有八次是设备没收到 CONNACK 就急着发 SUBSCRIBENetty 侧需要确保channelRead按序处理。5. 出包与排错JDK8 场景下的 NoClassDefFoundError 和优雅停机最后把项目打成可运行 zip 时类路径是最容易翻车的地方。5.1 NoClassDefFoundError: io/netty/util/timer如果你在 Spring Boot 项目里看到nested exception is java.lang.NoClassDefFoundError: io/netty/util/timer大概率是类路径里混入了 Netty 3 的包。Netty 3 的坐标是io.netty:nettyNetty 4 的坐标是io.netty:netty-all或io.netty:netty-common两者都有io.netty.util.Timer但包结构不同。Netty 4 把这个类移动到netty-common的io.netty.util包下而 NoClassDefFoundError 是运行时才出现的引用缺失编译期根本察觉不到。排查时先看依赖树mvn dependency:tree -Dincludesio.netty如果顶层带io.netty:netty:3.x就用 exclusions 排除dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId exclusions exclusion groupIdio.netty/groupId artifactIdnetty/artifactId /exclusion /exclusions /dependency这个错误经常出现在同时集成第三方 SDK 的项目里不是 Spring Boot 自身的问题。5.2 优雅停机的最后一道验证停服务时先关serverChannel再shutdownGracefully最后检查isRunning状态。shutdownGracefully会等待已连接 Channel 上的任务执行完默认有超时时间如果压测时发现停机要 30 秒以上就要考虑是自己 handler 里有阻塞操作而不是期待 shutdown 立即完成。验证停机和端口释放可以直接用ss -lntp | grep 1883该命令输出空行说明 Netty 的 EventLoopGroup 已经真正退出Spring Boot 容器随后才能安全关闭。最后再强调一次日志里如果出现leak detection警告说明 handler 里手动创建的 ByteBuf 没释放优先用Unpooled.copiedBuffer结合ReferenceCountUtil.release排查而不是先逃避内存泄漏检测。本文还有配套的精品资源点击获取
觉得有用,分享给同行:

为您的企业打造数字门面

稳重轻奢商务风格,端正雅致视觉,长效耐看不易过时。

立即咨询 →