
简介基于 Java AIO 开发的低延迟、高性能百万级 MQTT 客户端组件与 Broker 服务完整支持 MQTT v3.1、v3.1.1 和 v5.0 三套协议也支持 WebSocket MQTT 子协议兼容 mqtt.js、HTTP Rest API、遗嘱消息、保留消息、自定义消息转发与 Redis Pub/Sub 集群同时适配 GraalVM 原生编译和 Spring Boot 快速接入适合物联网、边缘计算等需要高并发低时延消息通信的开发者与架构师。压缩包内共 282 个文件以 221 个 Java 源码文件为主配合 15 个 Markdown 文档、13 个 XML、11 个 YAML 配置文件以及少量 JSON、脚本和 JAR 示例整体约 502KB目录划分清晰结构紧凑便于查阅。目前已有 280 人学习下载。透过源码可以了解 AIO 模式下 MQTT 编解码、会话管理、集群同步等核心实现并附带了 MQTT 客户端连接阿里云 Demo、HTTP 接口文件、Prometheus 与 Grafana 监控接入示例以及 Spring Boot 快速接入相关配置便于二次开发和生产落地适合想要深入研究高性能 MQTT 消息中间件的中高级 Java 工程师无论是学习源码还是工程选型都具备较高参考价值。1. 从“Java 到底能不能扛百万连接”说起AIO 选了哪条路车联网平台的设备接入层被问最多的问题就一个Java 到底能不能扛住百万级 MQTT 长连接拿 Arduino 这类 MCU 做采集下发的团队通常丢一个轻量客户端就完事但如果你负责的是自研物联网平台或 SCADA 汇聚层要在同一台 Broker 上同时持有几十万甚至一百万条设备连接选型就绕不开 AIO 和 NIO/Netty 之争。直接给结论Java AIO 的回调模型在纯长连接场景里比 Netty 更直白内存也可以规划得更紧凑作为百万级 MQTT 客户端组件和 Broker 服务的地基完全够用代价是要折腾一批 JDK 层面的细节。这篇笔记就从 AIO 线程模型一路拆到 MQTT 协议栈、客户端组件与 Broker 的落地路径和压测方法把血泪经验一并写上。2. AIO 线程模型与 MQTT 报文解码完成回调和无锁队列怎么协同2.1 AIO 与 NIO 的本质差别就绪通知与完成通知NIO 与 AIO 的区别是 java 面试题里常驻的高频考点但真到线上要你扛百万级连接时这道八股会变成非常具体的选型依据。NIO 的模型是“就绪通知”selector 告诉你有连接可读、有通道可写然后你的线程要自己去执行 read 或 write 系统调用。数据从内核缓冲区拷到用户态这一步耗时实打实发生在你的业务线程里。AIO 的模型是“完成通知”你提交一个 read 请求操作系统读完数据后把填好的 ByteBuffer 连同结果一起交到 CompletionHandler 的回调里。对业务代码来说省掉了“轮询状态 手动读数据”这个循环。这里有一个绕不开的实现事实Java AIO 在 Windows 上基于 IOCP是真正的内核异步在 Linux 上JDK 是用 epoll 事件循环加一层包装来模拟异步完成通知。也就是说你的工程落地不能指望内核原生 AIO但不影响使用 AIO 的编程模型。真正决定性能的反而是 AsynchronousChannelGroup 的线程池配置和回调代码的执行路径——这两个点后面专门讲。AIO 在纯长连接场景下的优势被低估了。百万连接里 99% 的消息是心跳和低频上行数据通道大多数时候是静的。NIO 模型下你得维护一个 selector 循环还要小心空转导致的 CPU 飙升AIO 模型下连接空闲时就是一个挂在 channel 上的 read 回调没有循环在跑。代码结构上也更直观一个连接的所有逻辑都写在一个回调里不需要在多个 handler 之间传状态机。2.2 MQTT 固定头与剩余长度在 AIO 回调里切半包不产生垃圾MQTT 报文结构很紧凑第一个字节的高四位是报文类型低四位是 DUP、QoS、RETAIN 标志紧跟着是剩余长度字段最长 4 字节每字节贡献 7 位有效数据最高位是延续标记最后是可变头和有效载荷。理解这个变长编码是写解码器的前提因为 TCP 是流协议回调里拿到的数据可能是一条报文、半条报文也可能粘了两条报文。在 AIO 回调里做协议解析核心原则是每个连接持有一个独立 ByteBuffer解析完半包要能回滚位置等下一次 read 回调攒够再继续。public class MqttDecoder { // 每个连接持有独立 buffer不跨连接共享 private final ByteBuffer buf ByteBuffer.allocate(8192); private final AsynchronousSocketChannel ch; public void readLoop() { buf.clear(); ch.read(buf, null, new CompletionHandlerInteger, Void() { Override public void completed(Integer result, Void attachment) { if (result 0) { // 对端关闭走连接清理逻辑 return; } buf.flip(); // 循环切包直到缓冲区里攒不够一条完整报文 while (true) { if (buf.remaining() 2) break; buf.mark(); // 记住报文起点半包时回滚 byte header buf.get(); int remaining decodeRemainingLength(buf); if (remaining 0 || buf.remaining() remaining) { // 半包剩余长度字段不完整或负载不足回滚等下次 buf.reset(); break; } byte[] payload new byte[remaining]; buf.get(payload); dispatch(header, payload); } buf.compact(); // 把剩余数据移到头部腾出空间 readLoop(); } Override public void failed(Throwable exc, Void attachment) { // 读失败通常意味着连接异常交给上层重连 } }); } // MQTT 剩余长度最多 4 字节每字节 7 位有效数据 private int decodeRemainingLength(ByteBuffer buf) { int multiplier 1; int value 0; for (int i 0; i 4; i) { if (!buf.hasRemaining()) return -1; byte b buf.get(); value (b 0x7F) * multiplier; if ((b 0x80) 0) { return value; } multiplier * 128; } return -1; // 超过 4 字节说明协议流有问题 } }这个解码器值得解释几个点。第一用 mark/reset 而不是把 position 手动赋值能避免多线程下 position 被意外修改的隐患虽然这里回调是串行的但养成用 mark/reset 的习惯更好。第二半包后 compact 会把未消费的字节移到缓冲区头部同时 position 指向数据末尾下一次 read 直接追加写入不用新建 buffer。第三payload 直接 new byte[remaining]如果这条连接 QoS2 大消息很多会产生频繁的字节数组分配后续优化可以用池化的 payload 容器但项目起步阶段简单拷贝更直观。2.3 回调线程不能干重活把业务丢出去把写操作排队AIO 的回调线程属于 AsynchronousChannelGroup。如果你在回调里做 JSON 解析、数据库鉴权、甚至阻塞的网络调用这个线程就被占住了同组其他连接的读写回调全部排队等候。百万连接下只需要几个回调线程被拖住几十毫秒吞吐就肉眼可见地往下掉。我一般会在回调里只做协议切包和消息转发业务逻辑一律丢进独立的业务线程池。final AtomicBoolean writePending new AtomicBoolean(false); final QueueByteBuffer writeQueue new ConcurrentLinkedQueue(); public void send(ByteBuffer msg) { writeQueue.offer(msg); // 只在没有写操作进行时才发起写避免 WritePendingException if (writePending.compareAndSet(false, true)) { doWrite(); } } private void doWrite() { ByteBuffer msg writeQueue.poll(); if (msg null) { writePending.set(false); return; } channel.write(msg, null, new CompletionHandlerInteger, Void() { Override public void completed(Integer n, Void a) { if (msg.hasRemaining()) { // 一次没写完继续写剩余字节 channel.write(msg, null, this); return; } doWrite(); // 写完整条消息后从队列取下一个 } Override public void failed(Throwable t, Void a) { // 连接异常标记连接失效并触发重连 } }); }这段代码解决的是 AIO 协议栈最常见的异常WritePendingException。一个 AsynchronousSocketChannel 上同一时间只允许一个写操作多个业务线程同时调 channel.write 会直接抛异常。加了一个 AtomicBoolean 标记写状态所有写请求先入队再由链式回调逐条发送。写队列的容量要设上限否则下游消费慢时客户端内存照样被堆爆这个在背压章节展开。3. 百万级 MQTT Client 组件连接管理、心跳抖动与背压控制3.1 AIO 客户端连接的最小闭环connect、CONNECT、CONNACK客户端组件的核心不是 TCP 连接建立而是连接建立后那一套 MQTT 会话握手。AIO 的 connect 调用是非阻塞的回调成功后才表示 TCP 层连通此时还要立即发送 CONNECT 报文。CONNECT 报文的三个关键参数是 clientId、keepAlive 和 cleanSession这三个参数直接决定连接被 Broker 接受后如何管理会话。public class MqttAioClient { private final AsynchronousSocketChannel channel; private final String clientId; private final int keepAlive; // 单位秒 public void connect(String host, int port) throws Exception { AsynchronousSocketChannel ch AsynchronousSocketChannel.open(); ch.connect(new InetSocketAddress(host, port), null, new CompletionHandlerVoid, Void() { Override public void completed(Void result, Void attachment) { // 连接成功后立刻组 CONNECT 报文 ByteBuffer packet MqttPackets.connectPacket(clientId, keepAlive); ch.write(packet, null, new CompletionHandlerInteger, Void() { Override public void completed(Integer n, Void a) { // 启动读循环等待 CONNACK 确认 readLoop(ch); } Override public void failed(Throwable t, Void a) { // 写失败按重连策略处理 } }); } Override public void failed(Throwable exc, Void attachment) { // 记录重连次数交给上层指数退避 } }); } }MqttPackets 是自定义的报文组装工具不是 JDK API这里要提醒读者Java 标准库不提供 MQTT 协议实现报文组装的位运算要自己写。connectPacket 内部需要按协议顺序填充固定头、剩余长度、协议名 MQTT、协议级别 4、连接标志位、keepAlive 双字节、clientId 长度和内容。这块建议单独封装成工具类因为客户端和 Broker 两侧都要复用。连接状态管理上客户端要维护一个状态字段CONNECTING、CONNECTED、RECONNECTING、CLOSED。CONNACK 报文回来后如果返回码为 0 才置为 CONNECTED其他返回码要能区分“服务端不可用”和“clientId 冲突”等。否则客户端会在 Broker 明确拒绝的情况下反复重连把服务端日志刷爆。3.2 百万客户端的重连策略指数退避和抖动窗口不能省连接数到十万级以后最怕的场景不是单个设备掉线而是成片设备同时掉线再同时重连。比如某个区域的 4G 基站抖动几万台设备同时断网网络恢复后全部涌向 Broker一瞬间的 connect 风暴能把服务端打挂。处理办法是双重抖动心跳时间抖动加退避抖动。ScheduledExecutorService scheduler Executors.newScheduledThreadPool(1); long keepAlive 300; // 秒 for (int i 0; i totalClients; i) { // 首跳时间均匀散落在 0~keepAlive 秒之间避免整点风暴 long jitter ThreadLocalRandom.current().nextLong(0, keepAlive * 1000L); scheduler.scheduleAtFixedRate(() - { if (!connected) { reconnectWithBackoff(); return; } channel.write(MqttPackets.pingReq(), null, writeHandler); }, jitter, keepAlive * 1000L, TimeUnit.MILLISECONDS); }心跳间隔设置对运营成本影响很大。常见的错误是把 keepAlive 设成 30 秒甚至 15 秒理由是“及时感知设备离线”。代价是三百万连接每秒产生 10 万条 PINGREQ整个 Broker 的 CPU 大部分耗在处理心跳上。我服务的物联网项目里车联网场景常用 300 秒或 600 秒配合 Broker 端 1.5 倍超时容忍设备掉线的感知延迟在一两分钟级别完全够用。真正要快速感知离线应该靠 MQTT 遗嘱消息在断线瞬间通知服务端而不是让心跳死快。重连退避公式业界常用的是 exponential backoff 加随机抖动delay min(maxDelay, base * 2^attempt) random(0, jitter)。base 取 1000msattempt 从 0 开始maxDelay 设上限 60 秒。没有随机抖动时同一批设备会按相同的时间序列重试问题只是从“同时重连”变成“每隔几秒同时重连”。加一个随机偏移后重连请求在整个时间轴上被摊开。3.3 背压控制Semaphore 限制在途消息数客户端发 PUBLISH 时TCP 层只保证字节发出去不保证对端应用层处理完。如果业务线程一股脑往 send 方法里灌消息而 Broker 或订阅方消费慢消息就会在客户端堆成山。QoS 1/2 场景下未确认消息堆积到一定量必须让上游感知压力。class MqttPublisher { // 信号量限制在途 QoS 1/2 消息数超过阈值阻塞调用方 private final Semaphore inflight new Semaphore(1024); private final AsynchronousSocketChannel channel; public void publish(String topic, byte[] payload) throws InterruptedException { inflight.acquire(); // 拿不到许可说明发送太快主动踩刹车 ByteBuffer msg MqttPackets.publishPacket(topic, payload); channel.write(msg, null, new CompletionHandlerInteger, Void() { Override public void completed(Integer n, Void a) { // 注意这里是写完成不是业务确认 inflight.release(); } Override public void failed(Throwable t, Void a) { inflight.release(); } }); } }Semaphore 初始值要按单连接的带宽延迟积估算。一个 QoS 1 消息平均 200 字节通道写缓冲常见 64KB信号量放 1024 意味着最多约 200KB 在途数据对 10Mbps 上行链路来说是合理的背压水位。如果用的是 QoS 2消息要等 PUBREL 完成才算一个完整的业务周期释放信号的时机要挪到 PUBCOMP 回调里。这里再强调一次上面的代码为了简洁在写回调里 release生产环境建议区分“写入完成”和“协议确认完成”否则流量大时问题会被掩盖。4. Broker 端做到百万连接线程池分片、订阅树与 QoS 状态机4.1 accept 循环与连接分片一个 AsynchronousChannelGroup 不够用自研 Broker 服务时入口是 AsynchronousServerSocketChannel 的 accept 回调每条连接建立后创建自己的 AsynchronousSocketChannel、读缓冲和写队列。但百万连接不能全部挂在同一个 AsynchronousChannelGroup 上。默认 group 的线程数和锁竞争在连接数过万后就会成为瓶颈回调线程被不同连接的事件互相挤兑。常见做法是分片创建多个 AsynchronousChannelGroup每个组配置 4 到 8 个线程把新连接按 clientId 或连接对象的 hash 散列到不同 group。这样一条连接的读写回调始终在同一组线程里执行各组的负载相对独立还能把锁竞争摊开。线程数不要无脑加IO 回调线程是 CPU 密集的超过核数以后上下文切换开销反而吃掉性能。分片之后Broker 的整体结构变成三层accept 层负责接入新连接分片层负责各自通道的读写回调业务层处理消息路由和持久化。accept 层的回调要非常轻量只做“创建连接对象 散列到分片 发起首次读”任何重操作放到业务层。否则 accept 回调慢了新连接建立速度直接受影响。4.2 订阅关系表通配符匹配与内存取舍Broker 的订阅路由是消息转发效率的关键。MQTT 主题按 / 分层订阅支持 和 # 通配符用平坦的 Map 存订阅关系没法高效匹配。我这边采用典型的 TopicTree 结构主题的每一层是树的一个节点节点挂订阅者集合。发布消息时按主题逐层向下找沿途收集 匹配的订阅者遇到 # 通配节点则把整棵子树都纳入。public class TopicNode { private final MapString, TopicNode children new ConcurrentHashMap(); // 当前节点下的订阅者读多写少用 CopyOnWriteArraySet 避免并发修改 private final SetSession subscribers new CopyOnWriteArraySet(); public void addSubscriber(String topic, Session session) { String[] tokens topic.split(/); TopicNode cur this; for (String token : tokens) { cur cur.children.computeIfAbsent(token, k - new TopicNode()); } cur.subscribers.add(session); } public void collect(String topic, ListSession out) { collect(topic.split(/), 0, out); } private void collect(String[] tokens, int idx, ListSession out) { if (idx tokens.length) { out.addAll(subscribers); return; } String token tokens[idx]; TopicNode child children.get(token); if (child ! null) { child.collect(tokens, idx 1, out); } // 匹配当前层的任意一个 token TopicNode any children.get(); if (any ! null) { any.collect(tokens, idx 1, out); } // # 匹配剩下的所有层收集整棵子树 TopicNode hash children.get(#); if (hash ! null) { hash.collectAll(out); } } }这个结构的问题在于内存占用。百万连接每个设备平均订阅 5 个主题就有 500 万个订阅关系。TopicNode、ConcurrentHashMap、CopyOnWriteArraySet 里每个对象都有不小开销粗估占用 800MB 到 1GB 堆内存是正常的。上线前必须按这个量级做内存预算不能只看连接本身。工业现场最常见的链路是 485 串口设备 - DTU - MQTT Broker - SCADA设备主题通常带设备 ID比如 devices/{deviceId}/telemetry这种高度规律的主题会把树的中间节点压得非常深内存估算时要按实际主题层级算。4.3 QoS 1/2 的会话状态机PUBREL 和 packetId 幂等Broker 对 QoS 1 的处理是收到 PUBLISH 后回复 PUBACK对 QoS 2 要跑完 PUBLISH - PUBREC - PUBREL - PUBCOMP 四次握手。QoS 2 的难点是幂等客户端因为超时会重传 DUP1 的 PUBLISHBroker 收到重复报文时必须能识别出来否则业务数据会被重复投递。class Qos2Session { // 保存已收到 PUBLISH 但还没收到 PUBREL 的 packetId private final MapInteger, MqttPublish inflight new ConcurrentHashMap(); public void onPublish(MqttPublish pub) { if (inflight.containsKey(pub.packetId())) { // 重复的 PUBLISH不投递业务直接回 PUBREC sendPubrec(pub.packetId()); return; } inflight.put(pub.packetId(), pub); deliverToSubscribers(pub); // 投递给订阅者 sendPubrec(pub.packetId()); } public void onPubrel(int packetId) { inflight.remove(packetId); sendPubcomp(packetId); } }这段逻辑背后是一个约定QoS 2 的“恰好一次”不是靠网络保证而是靠双方维护 packetId 状态实现去重。Broker 侧收到 PUBREC 前要把消息和相关状态保存好如果进程崩溃后消息还在但内存状态丢了重启后客户端重传的 PUBLISH 会被当成新消息再次投递破坏语义。cleanSessionfalse 的会话状态、QoS 2 的 inflight 表都需要定期落盘或同步到外部存储。这也是纯内存 Broker 在可靠性场景下最容易被挑战的地方设计阶段就要明确持久化边界。5. 从压测到上线的避坑清单连接数被卡、内存被吞、CPU 空转5.1 连接数卡在 65535 上下不去文件描述符和内核参数现象压测脚本不断创建新连接连接数到了某个值附近就上不去了等一两分钟又恢复一些随后再次卡住。原因这是最容易被忽视的系统层问题。单个进程的文件描述符上限 ulimit -n默认只有 1024 或 65535系统全局 fs.file-max 也可能成为天花板。Broker 每接受一条 TCP 连接至少要占用一个文件描述符百万连接就需要百万级 fd。另外客户端压测机出站连接受 net.ipv4.ip_local_port_range 限制默认端口范围不够时新连接会失败。解决服务端 ulimit -n 调到 1048576使用 systemd 管理的进程还要在服务单元里加 LimitNOFILE1048576sysctl -w fs.file-max1000000net.core.somaxconn 按 accept 队列调大。压测端机器同样要调而且压测机单机端口耗尽很常见多台压测机分摊才能接近真实百万连接。5.2 Direct Memory 爆掉ByteBuffer 没有复用现象压测跑到一半进程报 OutOfMemoryError: Direct buffer memory有时连错误日志都没有直接被 OS 的 OOM killer 干掉。原因AIO 的读写必须操作 ByteBuffer最常见的问题是每次读都新分配一个 buffer而且用的是 ByteBuffer.allocateDirect。堆外内存的回收依赖 GC 触发 cleaner高并发下新建速度远超回收速度堆外用量就一路上涨。还有一个隐蔽点JDK 的 ByteBuffer 对象如果没释放引用DirectByteBuffer 就不会被回收漏一个 Buffer 就是几百 KB 泄漏。解决每个连接在建立时分配一块固定大小的 Direct Buffer整个连接生命周期里反复使用用完通过 mark/compact 复用杜绝新建。启动参数加 -XX:MaxDirectMemorySize 设置上限比如 512MB 或 1GB超过直接报错方便排查。监控上盯 JMX BufferPool 里的 memoryUsed趋势上升就要排查泄漏。5.3 回调线程被业务拖死CPU 高但吞吐为 0现象连接数正常、CPU 跑满但消息吞吐掉到零偶尔有几个请求响应特别慢。原因把协议解码、消息路由、甚至数据库写入都塞在 AIO 回调线程里。AsynchronousChannelGroup 的线程池很小任一回调阻塞同组所有通道的读事件都在排队形成雪崩。线程 dump 里会看到业务代码的栈帧出现在 IO 线程上这是最直接的证据。解决回调线程里只做字节解析和报文切包消息处理的业务逻辑全部提交给独立线程池。用有界队列接住业务任务队列满了就触发该连接的背压等待不能无限往内存里塞。这里有一个检查技巧给线程池起名字用 ThreadFactory 把 IO 线程命名为 aio-callback- 前缀业务线程命名为 biz-worker- 前缀出问题时看线程名就知道是哪儿卡了。5.4 心跳抖动引起的大规模误下线现象某次网络抖动后Broker 上离线事件成片出现紧接着一大批客户端同时发起重连日志里全是连接建立与断开。原因客户端 keepAlive 设得太短比如 20 秒或 30 秒而 MQTT 服务端超过 1.5 倍 keepAlive 没收到任何报文就判定离线。移动网络下 4G/5G 切换、弱网环境里偶尔丢几个心跳包很正常超时判定触发后每台设备都走同样逻辑表现上就是集体误下线。解决把 keepAlive 调到 300 秒以上心跳只是“连接保活”的兜底不作为设备在线状态的主依据。设备真实上下线用遗嘱消息通知遗嘱在 TCP 断开的瞬间由 Broker 触发比心跳判定快得多也准得多。Broker 端对超时判定做 1.5 到 2 倍容忍不要把网络抖动放大成一次事故。5.5 重连风暴一百万客户端把刚起活的 Broker 打挂现象Broker 因发布或升级重启恢复后几分钟内收到海量 connect 请求新进程的线程池和认证模块被冲垮再次宕机。原因客户端在连接断开后立即重连有的实现甚至每 1 到 2 秒重试一次。Broker 重启的瞬时空档期越长客户端重试次数越多恢复后所有客户端还按固定节奏重试等于所有重试请求在同一时间窗口内撞在一起。解决客户端必须实现指数退避加随机抖动并且把单次重连的最大间隔限制在 60 秒内。Broker 端在启动初期可以限制 accept 速率比如前 30 秒只放行 20% 的连接其余丢进排队队列。如果认证逻辑连的是外部数据库要给认证模块加限流防止瞬时流量打挂下游。6. 压测与调优用阶梯式建连脚本盯住 P99 和 GC 停顿压测不能只盯着平均延迟和总吞吐百万连接场景下真正的敌人是长尾P99 延迟漂移、GC 停顿、堆外内存曲线。建连压测要分阶梯比如每 30 秒增加 1000 连接记录每一阶梯的建连耗时、失败数和内存增量直到连接数达到目标。一次性并发几十万连接测出来的数据没有参考价值真实业务里设备是逐步上线的。压测场景关键指标参考值瞬时建连每秒可接受连接数500010000百万连接稳态JVM 堆内存占比视订阅关系数量预留 1GB 以上余量QoS 1 消息下发P99 端到端延迟内网场景 100ms 以内Broker 重启恢复重连风暴下连接失败率0压测期间要同时采集 JVM 指标用 jstat -gcutil 看老年代和 GC 频率用 jcmd Thread.print 抓线程状态观察有没有回调线程长时间停在业务代码上。堆外内存用 JMX 的 BufferPool 监控建立连接阶段堆外内存如果持续上升基本上就是 ByteBuffer 泄漏先处理再继续测。调优顺序也有讲究先调系统层再调 JVM 直接内存和 GC 参数最后才动应用层的线程数和队列长度。顺序反了调了半天应用层发现瓶颈在 ulimit那就白忙一场。最后说一个我自己的习惯给每类线程池统一命名。AIO 回调线程、业务线程、重连调度线程分开命名压测出问题看 jstack 就能直接判断是哪个环节卡住。曾经有一次线上翻车连接建立慢了排查才发现调度线程池用的是默认名字 Thread-xxx全部对不上号只能一个个查栈顶。这个习惯帮我省了不止一次排障时间。希望这篇笔记把该踩的坑都提前标出来了真到你需要做百万级 MQTT 接入时能少走弯路。本文还有配套的精品资源点击获取