基于Java AIO的高性能MQTT Broker,百万级IoT网关连接实践

发布时间:2026/10/9 17:22:00
基于Java AIO的高性能MQTT Broker,百万级IoT网关连接实践
简介基于Java AIO实现的低延迟、高性能MQTT通信组件与Broker服务面向物联网、边缘计算和消息中间件方向的中高级开发者可支撑百万级连接场景帮助快速搭建私有的消息发布订阅通道。压缩包共282个文件约502KB以221个Java源码文件为核心辅以Markdown说明文档、XML/YAML配置文件、HTTP测试文件、JSON示例、Shell脚本及工程辅助文件等结构清晰便于阅读源码、部署调试与二次扩展。功能上覆盖MQTT v3.1、v3.1.1与v5.0协议解析支持WebSocket子协议、REST API、遗嘱消息、保留消息以及基于自定义消息处理转发和Redis Pub/Sub的集群机制还提供PrometheusGrafana监控对接。附带Spring Boot快速接入示例和阿里云MQTT连接Demo支持GraalVM原生编译兼顾协议细节、集群扩展与可观测性设计。目前已有280人学习/下载适合需要独立实现或二次封装MQTT服务的Java工程师也是高并发网络编程与物联网通信开发的实践参考。1. 百万级 MQTT 连接为什么 Java AIO 比 Netty 更适合你的下一套 IoT 网关当设备量级到了百万很多团队第一反应是上 Netty但我见过不少项目在 Netty 的线程模型上翻车连接数堆上去了心跳一多CPU 空转在 Selector 轮询上。这套基于 Java AIO 的 mqtt client 与 broker 组件走的是一条更直接的路——由操作系统把 IO 完成事件回调进来而不是业务线程自己轮询。它支持 MQTT v3.1、v3.1.1、v5.0 三套协议带 WebSocket 子协议和 HTTP REST API既能当 client 用也能当 broker 部署。对 Spring Boot 生态的 IoT 平台团队来说它最实际的价值是不用自己在 MQTT 编解码和连接管理上重复造轮子同时保留 GraalVM 编译成本机程序的选项适合从网关到边缘计算节点的各类落地场景。2. 架构与文件骨架从 AIO 事件模型到三版本协议栈的对齐关系2.1 AIO 的回调模型为什么低活跃长连接选 AIO 而不是 NIOMQTT 是典型的低活跃长连接协议设备连着服务器但绝大多数时间只在一分钟甚至几分钟一次的频率上报心跳。连接数越高、单连接活跃度越低NIO 模型里那个Selector.select()轮询的开销就越刺眼——每次轮询都要扫描一遍所有注册的 channel哪怕其中 99% 的连接没有任何数据到达。我见过一个压测到二十万连接的项目CPU 被 select 空转吃掉了三成业务线程反而分不到时间片。AIO 的处理思路完全不一样。它基于AsynchronousServerSocketChannel和CompletionHandler由操作系统在读写真正完成之后再回调你的 handler 方法。业务线程不需要主动去问“有没有数据”内核完成了直接通知你。这样一来连接的空闲时间越长、连接总量越大AIO 相对 NIO 的优势就越明显。下面是两类模型在一个典型 IoT 场景下的对比对比项NIOAIO事件获取方式Selector 轮询就绪事件内核完成 IO 后直接回调连接空闲时开销每次 select 仍扫描 channel空闲连接几乎零开销线程模型少量 IO 线程 业务线程池回调线程 业务线程池十万级以下连接足够胜任优势不明显百万级低活跃连接轮询空转明显更贴合长连接场景选型时有边界连接数在几千这个量级NIO 反而更好写因为轮询开销可忽略真正过十万、百万且低活跃AIO 的回调驱动才值得付出更多调试成本。这个项目把线程模型直接绑定到 JDK 的AsynchronousChannelGroup上我一般会先用默认线程数跑压测时再根据 CPU 核数调ioThreads一次只动一个变量。2.2 根目录 9 个关键文件一份能直接照着走的阅读路线拿到源码包后别急着从第一个类开始读先看根目录这 9 个文件它们把项目的构建方式、接口形态和协议核心三块已经交代清楚了文件作用建议动作mvnw.cmd/maven-wrapper.jarMaven Wrapper锁定 Maven 版本用它代替本机 mvn.editorconfig跨 IDE 的格式统一保持原样即可.gitignoreGit 忽略规则提交代码前确认mica-mqtt-api.httpREST API 调试脚本用 IDEA HTTP Client 直接执行ISSUE_TEMPLATE社区问题模板提 issue 时照模板填MqttDecoder.javaMQTT 解码器协议解析核心优先级最高MqttEncoder.javaMQTT 编码器与 Decoder 成对读DefaultMessageSerializer.java默认消息序列化器集群转发的扩展点我习惯先看mica-mqtt-api.http因为它把 broker 对外开放的 REST 接口一次性列出来了不需要翻代码就能知道这个组件对外提供了哪些管理能力。接着看MqttDecoder.java和MqttEncoder.java这两个类直接决定协议解析是否健壮。最后再看DefaultMessageSerializer.java它决定了自定义消息转发时消息体怎么被还原——需要做集群方案的团队改的就是这个类。2.3 协议版本分发一个 CONNECT 包怎么区分 v3.1 / v3.1.1 / v5.0同时支持三个协议版本最怕的就是解析错乱。MQTT 客户端的 CONNECT 报文结构是固定的固定头之后是可变头可变头里依次是协议名MQTT和协议级别Protocol Level这个 level 就是版本分发的关键。v3.1 对应3v3.1.1 对应4v5.0 对应5。按这个思路去读解码器分发逻辑应该是这样的public MqttMessage decode(int protocolLevel, ByteBuf in) { switch (protocolLevel) { case 3: return new Mqtt311Message(in).decodeV31(); case 4: return new Mqtt311Message(in).decodeV311(); case 5: return new Mqtt5Message(in).decodeV5(); default: throw new MqttException(unsupported protocol level: protocolLevel); } }这里的关键是 v5.0 不能复用 v3.1.1 的解析链路。v5.0 引入了属性Properties、Reason Code、Topic Alias、消息过期时间等机制报文结构比 v3.1.1 复杂不少。如果老设备用 v3.1.1 连上来服务端按 v3.1.1 走如果客户端声明支持 v5.0就按 v5.0 解析。两套逻辑共用一个固定头解析入口但可变头和 payload 的解析必须分开实现。我在接入时踩过的教训是不要试图用一个通用结构体兼容两个版本后面改一个字段就会把另一个版本的解析搞坏。3. 本地跑通全流程mvnw 打包、broker 启动与 mqtt.js 冒烟3.1 mvnw.cmd 构建没有 Maven 也能打包但先确认 JDK这套组件带了 Maven Wrapper所以本机不需要单独安装 Mavenmvnw.cmd会自动下载约定版本的 Maven 再执行构建。第一步先确认 JDK 环境建议至少 JDK 8配合 Spring Boot 2.x 用 JDK 8 或 11 最稳。命令行执行./mvnw.cmd clean package -Dmaven.test.skiptrueclean清掉上次构建产物package完成打包-Dmaven.test.skiptrue跳过测试以加快本地验证。首次执行时 wrapper 会下载 Maven 分发包和一堆依赖耗时主要卡在中央仓库。国内网络下如果长时间停在下载阶段在~/.m2/settings.xml里加阿里云镜像即可这与普通 Maven 项目是同一套配置mirror idaliyun/id mirrorOfcentral/mirrorOf urlhttps://maven.aliyun.com/repository/public/url /mirror提示构建失败最常见的原因是 JDK 版本与项目要求的字节码版本不匹配先看报错里是UnsupportedClassVersionError还是编译错误前者直接换 JDK 大版本后者才需要查依赖冲突。3.2 server 模式启动端口、SSL、心跳与保留消息配置清单这套组件既可以以 client 方式连接外部 broker也可以以 server 方式启动一个 broker。本地联调时我习惯先把 server 模式跑起来配置上重点关注这几个参数mica: mqtt: server: enabled: true host: 0.0.0.0 port: 1883 websocket-port: 8083 heartbeat-timeout: 180 ssl-enabled: false retain-enabled: true参数默认值参考说明port1883MQTT over TCP 监听端口websocket-port8083WebSocket 子协议监听端口mqtt.js 走这里heartbeat-timeout180服务端判定连接超时的时间单位秒ssl-enabledfalse是否开启 TLS生产环境建议开retain-enabledtrue是否启用保留消息存储注意一点项目里不同版本的配置项名可能略有出入解压后以包内application.yml实际为准但语义和上面这张表一致。heartbeat-timeout是避坑重点它指的是“多久没收到任何数据就断开”而不是“心跳间隔多久校验一次”。客户端把 keepalive 配成 60 秒时服务端这个值至少留到 180 秒原因在第五章第一条展开。WebSocket 端口要单独开mqtt.js 这类浏览器端客户端的连接都走它TCP 1883 端口对浏览器不可达。3.3 mqtt.js 冒烟三行代码验证订阅发布与 QoS 链路broker 起来之后建议用 mqtt.js 做一次冒烟验证不管你是 Arduino 设备接入还是前端联调这第一步都能确认“订阅与发布消息”这条主链路通不通const mqtt require(mqtt); const client mqtt.connect(ws://127.0.0.1:8083/mqtt, { clientId: smoke- Date.now(), cleanSession: true, }); client.on(connect, () { client.subscribe(demo/topic, { qos: 1 }); client.publish(demo/topic, hello from mqtt.js, { qos: 1 }); }); client.on(message, (topic, payload) { console.log(topic, payload.toString()); client.end(); });连接地址必须是ws://加上 WebSocket 端口路径上的/mqtt不能漏这对应项目里的 WebSocket MQTT 子协议实现。clientId在 v3.1.1 规范里最长 23 字节mqtt.js 5.x 生成的默认 ID 可能超长broker 会在 CONNACK 阶段直接拒绝。真遇到“客户端连不上但服务端日志显示连接已建立”这种诡异现象先查 clientId 再查路径。这条冒烟脚本跑通后再往上加遗嘱消息、保留消息、自定义消息转发这些能力才有意义。4. 核心代码拆读MqttDecoder、MqttEncoder 与消息序列化怎么联动4.1 剩余长度解析变长编码的四个字节与 256MB 上限MQTT 报文里最容易被读错的一段是“剩余长度”。它不是固定字节数而是 1 到 4 个字节的变长编码每个字节低 7 位是有效数据最高位是“是否还有后续字节”的标记。正因如此读这段逻辑时很多人会忘记处理后续字节导致一个报文没读完、流里下一个报文直接错位。解码器的核心逻辑大致是这样的int remainingLength 0; int multiplier 1; int encodedByte; do { encodedByte buffer.readUnsignedByte(); remainingLength (encodedByte 0x7F) * multiplier; multiplier * 128; } while ((encodedByte 0x80) ! 0);参数说明multiplier每轮乘 128 而不是 256因为每字节只有 7 位有效位循环终止条件是最高位为 0。MQTT 协议规定剩余长度最大 4 字节对应最大 256MB实际值为 268435455所以这个do-while循环最多执行 4 次。我在自己写的 decoder 里还会补一个计数保护超过 4 次直接抛异常避免恶意报文把读索引推到越界。这个细节在公网部署时必须加不然恶意客户端发一段全0x80的字节流就能拖垮解码线程。4.2 MqttEncoder 的写回策略占位符与 setByte 的两段式编码器看似比解码器简单实际有个麻烦剩余长度字段在可变头之前而它的值取决于可变头和 payload 的总长度。如果先写完 payload 再回头填剩余长度就得先算好长度再一次性输出做起来不顺。常见的做法是两段式先在固定头后面写一个占位字节等 body 写完拿到真实长度后再用setByte回填。示意如下public void encode(MqttMessage message, ByteBuf out) { int start out.writerIndex(); out.writeByte(message.header()); // 固定头 int lengthIndex out.writerIndex(); out.writeByte(0); // 剩余长度占位 // 写可变头和 payload int bodyLength out.writerIndex() - start - 2; out.setByte(lengthIndex, encodeRemainingLength(bodyLength)); }这里的关键参数是encodeRemainingLength剩余长度小于 128 时一个字节就能装下大于等于 128 时需要按变长规则拆成多字节。setByte只覆盖索引处的单字节不会影响已经写入的内容。每编码完一条消息还要记得检查写出去的字节数是否超出缓冲区最大帧限制。我在自己的项目里习惯把“最大报文长度”设成 1MB防止某个客户端一次性发布超大 payload 把内存打爆。4.3 DefaultMessageSerializer集群转发前为什么要过一层序列化单机 broker 转发消息时直接内存引用传递就够了但要做集群消息要经过 Redis pub/sub 广播到其他节点这时不能只发原始 payload必须把 topic、QoS、retain 标志一起带过去否则接收节点无法还原消息语义。DefaultMessageSerializer就是这套转发消息的序列化入口public class DefaultMessageSerializer implements IMessageSerializer { Override public byte[] serialize(String topic, byte[] payload, MqttQoS qos, boolean retain) { ByteBuffer buf ByteBuffer.allocate(payload.length 8); buf.put((byte) (qos.value() 1 | (retain ? 1 : 0))); buf.putShort((short) topic.length()); buf.put(topic.getBytes(StandardCharsets.UTF_8)); buf.put(payload); return buf.array(); } }这种结构把标志位放在最前面接收端反序列化时先读一个字节就能还原 QoS 和 retain。集群里每个 broker 节点都订阅同一个 Redis channel收到序列化消息后先看本地有没有对应 topic 的订阅者有才转发。我建议团队拿到这套组件后第一个自定义扩展就做在这里往消息体里加一个“来源节点 ID”字段这样排查消息是从哪个节点转发过来时就不用在 Redis 日志里大海捞针。5. 避坑指南百万连接场景下我踩过的 5 个 MQTT/AIO 的坑5.1 大量客户端连上就断heartbeat-timeout 不是心跳间隔现象设备端按 60 秒发一次心跳服务端heartbeat-timeout也配成 60 秒运行一晚上后掉了 10% 的连接重启后恢复循环往复。原因这个参数表示“服务端最多容忍多少秒没收到数据”而移动网络下设备可能因为信号问题延迟几分钟才发心跳客户端的心跳间隔和服务端的超时窗口之间没有任何容错空间。GC 停顿、线程调度延迟和 TCP 半开连接叠加后实际间隔经常超过 60 秒。解决把服务端超时设成客户端 keepalive 的 3 倍以上180 秒起步同时在服务端开启 TCP 半开探测定期发送探测包清理死连接。从那以后我再也不信“客户端配多少服务端就配多少”的说法全都按 3 倍窗口来。5.2 正常断开却触发遗嘱DISCONNECT 处理顺序的错位现象设备端按正常流程调用disconnect()主动断开但 broker 上遗嘱消息还是被发布了导致业务侧误报“设备离线”。原因连接关闭回调里broker 先执行了遗嘱发布逻辑然后再判断是否为优雅断开。很多实现把“连接异常终止”和“收到 DISCONNECT 后断开”混在了同一个清理链路里。v5.0 的 DISCONNECT 还带 Reason Code分支逻辑更多更容易写乱。解决在解析到 DISCONNECT 报文时打一个“优雅断开”标记连接关闭回调里先查标记标记存在就不发布遗嘱只清理会话只有标记不存在时才按异常断开处理并发布遗嘱。5.3 保留消息在集群节点丢失 QoS转发时别丢标志位现象单机下发布 retain 消息一切正常集群部署后有的订阅者拿到的是 QoS 0 而不是发布的 QoS 1甚至新订阅者取不到最后一条保留消息。原因消息经过 Redis pub/sub 转发时只传了 payload 字节没有把 retain 标志和 QoS 等级一并序列化。接收节点拿到裸字节后不知道这条消息是普通转发还是需要保留的 retained 消息只能按默认的 QoS 0 处理。解决转发消息统一走DefaultMessageSerializer把 QoS 和 retain 标志编码进消息头接收端反序列化后先恢复标志位再走保留消息的存储逻辑。序列化结构越早设计越好部署到生产再改要停机重发保留消息代价就大了。5.4 压测百万连接时堆外内存飙涨ByteBuffer 没有归还现象压测跑到二十万连接左右进程堆外内存持续上升最终触发 OOM 被杀堆内内存看起来没问题但操作系统层面的内存占用一直在涨。原因AIO 读回调里每次创建ByteBuffer或者缓冲区用完后没有执行release。JDK 的垃圾回收管不到堆外内存一旦池化失效泄漏是不可见的只能看着系统内存慢慢被吃掉。这是 AIO 编程最典型的坑也是大家说它“玄学”的主要原因。解决全程走池化ByteBuffer在CompletionHandler.completed()和failed()回调里都保证缓冲区归还Override public void completed(Integer result, Attachment attachment) { try { // 处理读到的消息 } finally { buffer.release(); // 归还到池不能漏 } }注意release()要放在finally里不能只放在正常分支末尾。读回调里某条消息解析抛异常缓冲区一样要归还否则一次异常就是一块堆外内存泄漏。5.5 生产环境别在 Windows 上压测IOCP 线程模型差异明显现象同一套源码在 Linux 上跑五十万连接没问题同事在 Windows 上压到一万连接 CPU 就满载甚至启动后直接卡死日志里没有任何异常。原因Windows 上 Java AIO 的底层是 IOCP线程挂起与完成端口的绑定机制和 Linux 下的 epoll 差异很大JDK 在 Windows 平台上的 AIO 线程膨胀行为也更激进。功能上能跑性能表现完全是两回事。解决生产环境锁 Linux x86_64开发环境只做功能验证性能数据一律以 Linux 压测结果为准。拿到这套组件后也别在 Windows 上折腾压测省下来的时间不如先把集群转发链路配好。6. 集群与监控的进阶验证从 Redis pub/sub 到 Prometheus 指标自检6.1 集群互通验证两个节点的 Redis pub/sub 联动起两个 broker 节点一个监听8083一个监听8084都接入同一个 Redis channel。用 mqtt.js 分别连两个节点节点 A 的客户端订阅节点 B 的客户端发布能收到就说明跨节点转发链路通了const a mqtt.connect(ws://127.0.0.1:8083/mqtt); const b mqtt.connect(ws://127.0.0.1:8084/mqtt); a.on(connect, () a.subscribe(cluster/topic)); b.on(connect, () b.publish(cluster/topic, hello cluster));6.2 Prometheus Grafana 指标位项目原生支持 Prometheus Grafana启动后直接拉 metric 端点即可。我落地时最关注的指标有三个当前连接数、每秒消息速率、心跳超时触发的断开次数。连接数曲线看容量水位消息速率对比压测脚本的发布速率可以定位 broker 的吞吐瓶颈心跳超时次数则能直接反映 5.1 那种配置错误。Grafana 面板我习惯把这三条曲线叠在一张图里压测时一眼能看出是配置问题还是资源问题。6.3 1000 到百万连接一个可持续复用的压测自检习惯每次改动连接池参数或编解码逻辑我都按这个顺序走一遍先 1000 连接跑通功能再看 10 万连接下的时延和内存曲线最后压到目标连接数观察 30 分钟长尾。重点不是冲高点而是看 Grafana 里堆外内存和线程数曲线是不是平稳。实现上优先用现成的 mqtt 压测工具没有合适工具时用 mqtt.js 起多个进程模拟设备也够用。从那以后我每接手一个 MQTT broker 资源都会强制走一遍这个流程本地 mqtt.js 冒烟 → 集群 Redis 联动 → Prometheus 曲线压测自检。这套源码包解压后建议你也按这个顺序走一遍重点核对heartbeat-timeout、遗嘱处理分支和DefaultMessageSerializer里的标志位设计。希望帮到你。本文还有配套的精品资源点击获取