简介:基于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的集群机制,还提供Prometheus+Grafana监控对接。附带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 而不是 NIO
MQTT 是典型的低活跃长连接协议:设备连着服务器,但绝大多数时间只在一分钟甚至几分钟一次的频率上报心跳。连接数越高、单连接活跃度越低,NIO 模型里那个Selector.select()轮询的开销就越刺眼——每次轮询都要扫描一遍所有注册的 channel,哪怕其中 99% 的连接没有任何数据到达。我见过一个压测到二十万连接的项目,CPU 被 select 空转吃掉了三成,业务线程反而分不到时间片。
AIO 的处理思路完全不一样。它基于AsynchronousServerSocketChannel和CompletionHandler,由操作系统在读写真正完成之后再回调你的 handler 方法。业务线程不需要主动去问“有没有数据”,内核完成了直接通知你。这样一来,连接的空闲时间越长、连接总量越大,AIO 相对 NIO 的优势就越明显。下面是两类模型在一个典型 IoT 场景下的对比:
| 对比项 | NIO | AIO |
|---|---|---|
| 事件获取方式 | Selector 轮询就绪事件 | 内核完成 IO 后直接回调 |
| 连接空闲时开销 | 每次 select 仍扫描 channel | 空闲连接几乎零开销 |
| 线程模型 | 少量 IO 线程 + 业务线程池 | 回调线程 + 业务线程池 |
| 十万级以下连接 | 足够胜任 | 优势不明显 |
| 百万级低活跃连接 | 轮询空转明显 | 更贴合长连接场景 |
选型时有边界:连接数在几千这个量级,NIO 反而更好写,因为轮询开销可忽略;真正过十万、百万且低活跃,AIO 的回调驱动才值得付出更多调试成本。这个项目把线程模型直接绑定到 JDK 的AsynchronousChannelGroup上,我一般会先用默认线程数跑,压测时再根据 CPU 核数调ioThreads,一次只动一个变量。
2.2 根目录 9 个关键文件:一份能直接照着走的阅读路线
拿到源码包后,别急着从第一个类开始读,先看根目录这 9 个文件,它们把项目的构建方式、接口形态和协议核心三块已经交代清楚了:
| 文件 | 作用 | 建议动作 |
|---|---|---|
mvnw.cmd/maven-wrapper.jar | Maven Wrapper,锁定 Maven 版本 | 用它代替本机 mvn |
.editorconfig | 跨 IDE 的格式统一 | 保持原样即可 |
.gitignore | Git 忽略规则 | 提交代码前确认 |
mica-mqtt-api.http | REST API 调试脚本 | 用 IDEA HTTP Client 直接执行 |
ISSUE_TEMPLATE | 社区问题模板 | 提 issue 时照模板填 |
MqttDecoder.java | MQTT 解码器 | 协议解析核心,优先级最高 |
MqttEncoder.java | MQTT 编码器 | 与 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 对应3,v3.1.1 对应4,v5.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,所以本机不需要单独安装 Maven,mvnw.cmd会自动下载约定版本的 Maven 再执行构建。第一步先确认 JDK 环境,建议至少 JDK 8,配合 Spring Boot 2.x 用 JDK 8 或 11 最稳。命令行执行:
./mvnw.cmd clean package -Dmaven.test.skip=trueclean清掉上次构建产物,package完成打包,-Dmaven.test.skip=true跳过测试以加快本地验证。首次执行时 wrapper 会下载 Maven 分发包和一堆依赖,耗时主要卡在中央仓库。国内网络下如果长时间停在下载阶段,在~/.m2/settings.xml里加阿里云镜像即可,这与普通 Maven 项目是同一套配置:
<mirror> <id>aliyun</id> <mirrorOf>central</mirrorOf> <url>https://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| 参数 | 默认值参考 | 说明 |
|---|---|---|
port | 1883 | MQTT over TCP 监听端口 |
websocket-port | 8083 | WebSocket 子协议监听端口,mqtt.js 走这里 |
heartbeat-timeout | 180 | 服务端判定连接超时的时间,单位秒 |
ssl-enabled | false | 是否开启 TLS,生产环境建议开 |
retain-enabled | true | 是否启用保留消息存储 |
注意一点:项目里不同版本的配置项名可能略有出入,解压后以包内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里的标志位设计。希望帮到你。
本文还有配套的精品资源,点击获取