资讯动态

基于Netty与MQTT 3.1.1自研物联网服务端和客户端实践

发布时间:2026/9/10 22:33:46 来源:尧图企业网站定制
简介基于Netty、MQTT 3.1.1、Spring Boot和JDK 8构建的MQTT服务端与客户端完整工程面向具备Java基础、希望深入物联网通信或准备毕业设计的开发者。工程代码结构清晰包含服务端、客户端、通用模块及启动器涵盖连接建立、消息发布订阅、会话保持等核心逻辑。压缩包共113个文件以87个Java源码为主辅以XML/YML配置、JKS安全证书、Spring factories及说明文档整体仅153KB轻量且便于阅读。目前已有536人学习下载。通过阅读源码与部署说明可以掌握Netty的异步事件驱动模型理解MQTT 3.1.1协议的报文交互细节并学习Spring Boot多模块项目的组织方式同时参考证书配置与测试用例为实际IoT场景或毕业设计提供可复用的实战模板。1. 用 netty mqtt 3.1.1 自研 MQTT 服务端与客户端到底解决什么问题做物联网接入的人大概率遇到过这个选择题设备端是 stm32、4G 模块这类跑不动 HTTP 的终端云侧要接收大量长连接还需要把每一条上行数据对应到具体业务订单上。第一反应是部署现成的 broker但接进去才发现设备鉴权要走自己的账号体系、上行数据要落库、下行指令要关联业务逻辑改 broker 的插件往往比写一套新服务还费劲。我一般直接选择用 netty 写。这套标题给定的组合——netty mqtt 3.1.1 springboot jdk8——是物联网服务端最常见的基线。mqtt 3.1.1 协议本身足够薄netty 官方提供了 codec真正需要自己写的只有连接管理、会话保持、主题路由和业务钩子。它能解决的是「在不改现成 broker 的情况下把接入层完全攥在自己手里」的问题适合做充电桩、传感器网关、工业采集这类既要接设备又要接业务系统的项目。下面从协议点开始逐层把这个服务端和客户端落地。2. MQTT 3.1.1 协议要点与 netty 编解码选型2.1 mqtt 3.1.1 报文体例控制包类型、剩余长度与 QoSmqtt 3.1.1 的报文分固定头、可变头、载荷三段。固定头第一个字节高 4 位是控制包类型低 4 位是标志位第二字节起是剩余长度表示后续可变头加载荷的总字节数用 1 到 4 个字节编码每个字节低 7 位是有效数据最高位是「是否还有后续字节」的标记。所以小于 128 的长度只用一个字节超过就得递增字节数服务端解码时必须按这个规则完整读出一个包再向上抛不能按 TCP 流边界切业务报文。下面是接入服务端前必须认全的报文类型报文类型含义方向使用场景CONNECT连接请求客户端到服务端设备首次接入携带 clientId、用户名、遗嘱CONNACK连接确认服务端到客户端服务端返回接受或拒绝码PUBLISH发布消息双向数据上行与指令下行PUBACKQoS 1 确认双向收到 QoS 1 消息后回执SUBSCRIBE订阅请求客户端到服务端设备订阅下行主题SUBACK订阅确认服务端到客户端逐主题返回授予的 QoSPINGREQ / PINGRESP心跳请求 / 响应双向保活DISCONNECT主动断开客户端到服务端正常下线不触发遗嘱项目标题锁在 mqtt 3.1.1 而不是 5.0是因为 3.1.1 是当前嵌入式侧兼容性最好的版本。stm32、4G 模组自带的 mqtt 协议栈大多基于 3.1.1netty 的 codec 对它支持也最完整。5.0 增加了属性、Reason Code、共享订阅等能力但很多设备端 SDK 还没跟上做接入层选 3.1.1 是求最省事的路径。2.2 为什么选 netty 而不直接用现成 broker现成 broker 适合「有现成接入、有标准主题规划」的项目比如快速把设备接到 EMQX 或 Mosquitto 上再通过 webhook 同步到业务库。但一旦涉及私有鉴权、设备影子、指令下发要拼接业务参数插件机制就变成限制。netty 做这件事的收益是三层连接层线程模型是现成的 NIO不需要自己管连接和线程的对应codec 包已经实现 mqtt 3.1.1 的编解码省掉最麻烦的报文解析业务 handler 可以直接写在 pipeline 里和 springboot 的 service 层无缝衔接。对比下来自研 netty 服务端的可接受代价是会话管理、遗嘱、QoS 重发逻辑都要自己写。一个折中思路是「broker 对外 netty 做接入网关」只把连接层握在自己手里转发给内部 broker。这种架构实践中很常见但复杂度更高。标题里既然给了 netty springboot jdk8 的组合我按完全自研的实现往下走对任何项目的参考价值都更大。2.3 引入 netty-codec-mqtt跑通最小服务端jdk8 对应 netty 4.1.x 没有任何兼容问题引入 codec 模块即可不需要引 netty-all。pom 里两个依赖就够dependency groupIdio.netty/groupId artifactIdnetty-codec-mqtt/artifactId version4.1.100.Final/version /dependency dependency groupIdio.netty/groupId artifactIdnetty-transport/artifactId version4.1.100.Final/version /dependency最小服务端只需要一个 ServerBootstrapEventLoopGroup boss new NioEventLoopGroup(1); EventLoopGroup worker new NioEventLoopGroup(); try { ServerBootstrap bootstrap new ServerBootstrap(); bootstrap.group(boss, worker) .channel(NioServerSocketChannel.class) .childHandler(new ChannelInitializerSocketChannel() { Override protected void initChannel(SocketChannel ch) { ch.pipeline().addLast(new MqttDecoder(1024 * 1024)); ch.pipeline().addLast(MqttEncoder.INSTANCE); ch.pipeline().addLast(new MqttServerHandler()); } }) .bind(1883).sync().channel().closeFuture().sync(); } finally { boss.shutdownGracefully(); worker.shutdownGracefully(); }boss 线程只负责 acceptworker 线程处理每个连接的读写这个拆分是 netty 高性能的基础。MqttDecoder 构造参数 maxBytesInMessage 是单个 mqtt 报文允许的最大字节数默认只有 8KB 左右不适用于图片或批量数据上报按业务负载放大到 1MB 比较常见。MqttEncoder 是无状态的用单例直接加入 pipeline。跑通这一步之后用任意 mqtt 客户端连上来不会报错但也没有响应因为业务 handler 还没写下一章补上。3. 服务端实现连接鉴权、心跳保持与主题路由3.1 CONNECT 报文处理从解析到鉴权再到遗嘱服务端 handler 继承 SimpleChannelInboundHandler第一个要处理的入站消息就是 CONNECT。MqttConnectMessage 的 payload 里能拿到 clientId、用户名、密码variableHeader 里能拿到 keepAlive、cleanSession 和遗嘱标记。处理顺序建议固定为先查 clientId 是否已被占用再验证用户名密码最后根据 cleanSession 决定是否恢复会话遗嘱消息则要先存到 Session 对象里再返回 CONNACK。Override protected void channelRead0(ChannelHandlerContext ctx, MqttMessage msg) { if (msg instanceof MqttConnectMessage) { MqttConnectMessage conn (MqttConnectMessage) msg; String clientId conn.payload().clientIdentifier(); String username conn.payload().userName(); String password new String(conn.payload().passwordInBytes(), StandardCharsets.UTF_8); boolean cleanSession conn.variableHeader().isCleanSession(); int keepAlive conn.variableHeader().keepAliveTimeSeconds(); if (SessionManager.isOnline(clientId)) { // mqtt 3.1.1 规范同 clientId 新连接接管旧连接先踢旧的 SessionManager.kick(clientId); } if (!authService.check(username, password)) { ctx.writeAndFlush(connAck(MqttConnectReturnCode.CONNECTION_REFUSED_BAD_USER_NAME_OR_PASSWORD)); ctx.close(); return; } ctx.writeAndFlush(connAck(MqttConnectReturnCode.CONNECTION_ACCEPTED)); Session session SessionManager.create(clientId, ctx.channel(), cleanSession, keepAlive); if (conn.variableHeader().isWillFlag()) { session.setWill(conn.payload().willTopic(), conn.payload().willMessageInBytes()); } } }CONNACK 构造时固定头类型用 CONNACK可变头要回传 sessionPresent 标志和返回码返回码用 MqttConnectReturnCode 枚举拒绝码要区分 0x04 用户名密码错误和 0x05 未授权方便设备侧按返回码决定是否重连。CONNECT 处理是在 netty 的 EventLoop 线程里执行的鉴权查 redis 没问题但如果走慢 SQL 查询建议把鉴权放到独立业务线程池里异步完成拿到结果后再 writeAndFlush 回 CONNACK否则会阻塞整个 EventLoop 上其他连接的数据读取。这是 netty 服务端最容易踩的并发陷阱。3.2 心跳保活IdleStateHandler 与 userEventTriggered 的配合mqtt 3.1.1 的心跳机制是客户端在 keepAlive 时间窗口内至少发一个 PINGREQ服务端回 PINGRESP。服务端要做的是在 CONNECT 里拿到 keepAlive 值然后给这条连接的 pipeline 加一个 IdleStateHandler超时事件会以 IdleStateEvent 形式触发 handler 的 userEventTriggered。这个回调在 netty 服务端里是面试高频问题实际使用难点在于阈值设多少。IdleStateHandler 参数含义服务端常见设置readerIdleTime读空闲超时超过该时长没收到任何数据触发事件keepAlive 的 1.5 倍writerIdleTime写空闲超时超过该时长未写出数据触发事件0表示不启用allIdleTime读或写任意一方向空闲即触发0表示不启用官方 keepAlive 只约束客户端必须发数据没有规定服务端容忍多久。把 readerIdleTime 精确设成 keepAlive 会导致一次网络抖动就误踢设备常见做法取 1.5 倍给重传和数据包延迟留余量。在 CONNECT 校验通过后加 idle handler踢人的逻辑集中在 userEventTriggered 里ch.pipeline().addLast(idle, new IdleStateHandler(keepAlive * 3 / 2, 0, 0)); Override public void userEventTriggered(ChannelHandlerContext ctx, Object evt) throws Exception { if (evt instanceof IdleStateEvent) { IdleStateEvent idle (IdleStateEvent) evt; if (idle.state() IdleState.READER_IDLE) { // 超时未收到客户端任何报文包括 PINGREQ SessionManager.closeByChannel(ctx.channel()); ctx.close(); return; } } super.userEventTriggered(ctx, evt); }注意别把读超时全局写死不同设备上报周期差别很大充电桩可能 60 秒心跳温湿度传感器可能 5 秒从 CONNECT 里动态取 keepAlive 再算阈值才是对的。另外 PINGREQ 本身就能刷新读空闲时间所以即使客户端不发业务数据只发心跳也不会被误判离线。3.3 订阅、发布与主题通配符匹配SUBSCRIBE 报文解析出主题和期望 QoS 后服务端把它记入订阅表用 ConcurrentHashMap 按订阅主题维度存放 set of Channel。mqtt 3.1.1 的主题层级用斜杠分隔通配符只有两个 代表单层任意/ 用于多层剩余部分且必须以单个主题过滤器结尾的方式出现。匹配时不能只做字符串前缀要看层级static boolean match(String filter, String topic) { if (filter.equals(topic)) { return true; } if (!filter.contains() !filter.contains(#)) { return false; } String[] f filter.split(/); String[] t topic.split(/); for (int i 0; i f.length; i) { if (i t.length) { return false; } if (.equals(f[i])) { continue; } if (#.equals(f[i])) { return true; } if (!f[i].equals(t[i])) { return false; } } return f.length t.length; }只能在过滤器最后一段代表「剩余所有层级」所以匹配到 # 直接返回 true。 必须匹配正好一层t 的层级数和 f 对齐。广播时遍历订阅表逐条匹配命中后构造 MqttPublishMessage 写入对应 channel注意 qos 取「发布 QoS 与订阅 QoS 的较小值」这是协议规定没必要让服务端强行升 QoS。写失败说明对端掉线从订阅表移除该 channel再按其 session 遗嘱处理。4. 客户端实现SpringBoot 集成、连接、订阅与重连4.1 springboot 连接参数配置与注入服务端写完之后需要一个能跑的客户端来验证实际项目里这个客户端通常是网关或后端服务。Spring Boot 2.7.x 是 jdk8 下的稳妥选择配置集中在 application.yml 里。把这些参数做成配置类的好处是换环境不用改代码设备侧 MQTT 连接参数和业务参数分开管理。mqtt: broker-url: tcp://localhost:1883 client-id: gateway-001 username: app password: app-secret keep-alive: 30 clean-session: true reconnect-interval-ms: 5000配置项说明broker-url格式为 tcp://host:portssl 场景用 ssl://client-id在 broker 内必须唯一重复会导致旧连接被接管keep-alive客户端保活周期服务端按 1.5 倍做读超时clean-session为 true 时服务端不保留离线消息与会话reconnect-interval-ms断线后首次重连延迟重试可做指数退避4.2 用 netty 客户端封装 connect、subscribe、publish客户端同样用 Bootstrap只是方向相反。pipeline 里加 codec 后连接建立不代表 mqtt 连接建立必须发 CONNECT 并收到 CONNACK 才算 mqtt 层就绪。订阅和发布都在 CONNACK 成功之后进行Component public class MqttClient { private Bootstrap bootstrap; private volatile Channel channel; public void connect(String host, int port) throws Exception { EventLoopGroup group new NioEventLoopGroup(1); bootstrap new Bootstrap(); bootstrap.group(group) .channel(NioSocketChannel.class) .handler(new ChannelInitializerSocketChannel() { Override protected void initChannel(SocketChannel ch) { ch.pipeline().addLast(MqttEncoder.INSTANCE); ch.pipeline().addLast(new MqttDecoder(1024 * 1024)); ch.pipeline().addLast(new SimpleChannelInboundHandlerMqttMessage() { Override protected void channelRead0(ChannelHandlerContext ctx, MqttMessage msg) { if (msg instanceof MqttConnAckMessage) { // 根据返回码决定是否订阅业务主题 } else if (msg instanceof MqttSubAckMessage) { // 订阅确认按 topic 检查返回的 QoS } else if (msg instanceof MqttPublishMessage) { // 下行指令转交业务线程处理 } } }); } }); ChannelFuture future bootstrap.connect(host, port).sync(); channel future.channel(); sendConnect(); } private void sendConnect() { MqttFixedHeader fixed new MqttFixedHeader(MqttMessageType.CONNECT, false, MqttQoS.AT_MOST_ONCE, false, 0); MqttConnectVariableHeader vh new MqttConnectVariableHeader( mqtt;clientjava, 4, true, true, true, 1, true, 30); MqttConnectPayload payload new MqttConnectPayload(clientId, null, null, username, password.getBytes()); channel.writeAndFlush(new MqttConnectMessage(fixed, vh, payload)); } }connection 和业务 handler 分开是保持可读性的关键。重连不能靠 catch 异常后立刻 new bootstrap常见做法是监听连接失败用 eventLoop 的 schedule 做延迟重连第一次 5 秒、后续翻倍到 60 秒封顶future.addListener((ChannelFutureListener) f - { if (!f.isSuccess()) { long next Math.min(reconnectIntervalMs * 2L, 60_000L); channel.eventLoop().schedule(() - connect(host, port), next, TimeUnit.MILLISECONDS); } });延迟任务要挂在 eventLoop 上而不是用 spring 的定时线程池这样重连和连接事件天然在同一线程内执行避免 Channel 并发写的问题。收到 CONNACK 且返回码不是 0 时要按码区分处理用户名密码错误直接停掉重连不要无限重试浪费 broker 资源。4.3 jdk8 的 CompletableFuture 解决连接异步链纯回调写法在「连接成功 → 订阅 → 再发布」三级联动时容易嵌套过深。jdk8 的 CompletableFuture 能把 netty 的回调改写成链式风格也是 jdk8 这个基线里值得用的新特性。把发送 CONNECT 改为返回 futurehandler 收到 CONNACK 时 complete 它CompletableFutureMqttConnAckMessage ackFuture new CompletableFuture(); channel.writeAndFlush(connectMessage); ackFuture.thenAccept(ack - { if (ack.variableHeader().connectReturnCode() MqttConnectReturnCode.CONNECTION_ACCEPTED) { subscribe(dev/ clientId /cmd, MqttQoS.AT_LEAST_ONCE); } });参数说明thenAccept 在 CONNACK 到达后执行订阅动作放在这里可以保证顺序complete 只能在连接关闭前调用一次再次重连需要 new 一个新的 future。lambda 表达式在这个链式调用里比匿名内部类省很多代码而且 jdk8 的 CompletableFuture 不依赖额外依赖包。5. 验证、压测与三个高频坑5.1 用 mosquitto 客户端端到端验证服务端和客户端都写完直接互相测不可靠要有一个独立第三方验证。本机装 mosquitto-clients 工具订阅和发布分开验证mosquitto_sub -h 127.0.0.1 -p 1883 -t dev//data -v mosquitto_pub -h 127.0.0.1 -p 1883 -t dev/001/data -m {t:26.5}订阅端能打印出dev/001/data {t:26.5}说明主题匹配和广播逻辑正常。再测遗嘱和踢人开两个订阅端杀掉其中一个观察另一个是否收到遗嘱消息。验证心跳时把 CONNECT 里的 keepAlive 改成 3 秒等待 5 秒看服务端是否按 1.5 倍阈值断开该连接。5.2 高频坑一nested exception is java.lang.NoClassDefFoundError: io/netty/util/timer这个报错在 jdk8 项目里几乎都能定位到同一个原因类路径里混入了 netty 3.x。io/netty/util/timer是 netty 3 的包路径netty 4 改成了io.netty.util代码在编译期引用到 4.x运行期被某个传递依赖把 3.x 拉进来就抛 NoClassDefFoundError。排查手段是mvn dependency:tree -Dincludesio.netty:netty*把 netty 3 的传递依赖用 exclusion 排除再用 dependencyManagement 把 netty 版本统一到 4.1.x。spring boot 2.7 自带的 netty 版本管理能兜住大部分场景直接标的 4.1.100.Final 要和 boot 管理版本对齐。5.3 高频坑二spring boot 版本太高与 jdk8 不兼容spring boot 3.x 强制要求 jdk17用 jdk8 直接启动会出现 class version 错误表现是几十个类加载异常容易被误判成 netty 问题。jdk8 基线对应 spring boot 2.7.x这个分支的依赖管理里 netty 版本也能正常工作。如果业务要求必须升级 boot 3.x先升 jdk 后升 boot升级过程里 netty 的 API 没有破坏性变化但消息处理里如果用了 jdk8 的 Optional、CompletableFuture 相关代码注意 boot 3 默认 web 容器和序列化行为有变化。5.4 高频坑三心跳误踢与 QoS 1 重复投递读空闲阈值设成与 keepAlive 一致是误踢的主要原因按 1.5 倍已经说过。QoS 1 的重复投递是另一个容易忽略的问题客户端发送 PUBLISH 后如果没收到 PUBACK重发同一条消息服务端会因为 messageId 不同而当作新消息投递业务侧必须对 payload 做幂等常见做法是以设备上报的 seq 字段做去重。调这种问题最直接的工具是 packet capture 抓包看报文流转。调试技巧在服务端 pipeline 最前面加一个打印十六进制报文的 handler用 ByteBufUtil.hexDump 输出原始字节能一眼看出客户端到底发的 CONNECT 还是 PINGREQ省去猜报文格式的时间。这个方法在验证自定义客户端时比看业务日志更快也是排查协议层问题最通用的收尾手段。本文还有配套的精品资源点击获取

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价