ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

Java AIO与MQTT百万级长连接实战:从线程模型到Broker调优全解析

2026/10/4 10:53:58 拓冰建站 浏览量
Java AIO与MQTT百万级长连接实战:从线程模型到Broker调优全解析 简介基于 Java 异步 IOAIO技术打造的高性能消息队列遥测传输MQTT客户端与服务端组件专为物联网、边缘计算及消息服务器开发者设计旨在解决海量设备接入与低延迟消息转发的工程难题。项目完整支持 MQTT v3.1、v3.1.1 与 v5.0 协议具备 WebSocket 子协议、REST 接口、遗嘱消息、保留消息等实用特性同时提供基于 Redis 发布订阅的集群扩展、Spring Boot 快速接入示例、阿里云 MQTT 连接演示并支持 Prometheus 与 Grafana 监控、GraalVM 原生编译以及自定义消息转发。压缩包共二百八十二个文件其中二百二十一个 Java 源文件构成核心逻辑辅以 XML 与 YAML 配置、Markdown 文档、脚本及 Maven 包装器整体仅五百零二 KB目录结构清晰便于深入研读编解码、消息流转、集群通信等关键模块。已有二百七十九人学习下载适合具备 Java 网络编程基础、希望掌握 MQTT 底层机制或搭建高并发物联网消息系统的开发者。1. 百万级 MQTT 连接为什么卡在 Java 的 BIO 模型上先抛一个反直觉的结论用 Java 写 MQTT Broker 和客户端真正逼你上 AIO异步 IO的不是 CPU而是线程。BIO 模型下一个连接一个线程8 核机器开到 2000 连接就开始频繁上下文切换GC 线程和 IO 线程互相抢时间片到 5000 连接基本就废了。而百万级连接意味着你不可能给每个连接配一个线程只能让少量线程复用把「等数据」这件事交给内核。这就是 java aio 在这个标题里的核心价值用少量 IO 线程承载海量连接把线程数从「连接数」降级为「CPU 核数」。很多人一听到百万级 MQTT 就本能地去调 JVM 堆内存实际上堆内存从来不是瓶颈。真正的问题是每一条 MQTT 连接都有自己的状态机、心跳定时器、待发送队列这些对象如果不做池化和复用百万连接光对象头就能吃掉几个 GB。所以做这个组件第一步不是写代码而是想清楚你的连接状态用什么数据结构存、心跳任务怎么调度、发送缓冲区怎么分配。这篇笔记我会按「客户端组件」和「broker 服务」两条线拆开讲把 AIO 的线程模型、回调写法、背压处理、参数调优和踩坑记录都摊开最后给你一套能在本地验证效果的完整方案目标是让你能区分「架构上能做百万级」和「上线后真的撑得住百万级」。2. 先看透 Java AIO 的线程模型为什么说它是为 MQTT 长连接量身定做的2.1 AIO 和 BIO、NIO 的本质区别回调还是轮询Java 的三种 IO 模型里BIO 是阻塞等数据每个连接一个线程代码好写但线程数是硬上限NIO 是同步非阻塞一个线程轮询多个 Channel 的就绪事件能扛的连接数上去了但你的业务逻辑得自己处理半包、粘包和状态机AIO 更进一步把「数据什么时候就绪」这个通知机制交给内核你注册一个 CompletionHandler读操作完成后系统直接回调你的方法。对应到 MQTT 协议上这个区别非常明显。MQTT 的报文是变长的有固定头、可变头、载荷三段你要自己处理「读到的字节不够一个完整报文」的边界情况。BIO 下你在 read() 里等写起来顺但线程废了NIO 下你要在 SelectionKey 里反复判断当前积累了多少字节AIO 下你发起一个 read 操作内核把数据填进 ByteBuffer 后回调你的 completed() 方法你在这里面做报文解析、状态流转然后再次发起 read形成链式回调。这里有个关键点AIO 的回调线程是内核的 IO 线程池里的线程不是你的业务线程。所以你在 completed() 里绝对不能做耗时操作比如查数据库、加锁、做复杂编解码否则会把 IO 线程池堵死其他连接的读写全部延迟。常见的做法是在回调里只做「拆包 把完整报文丢进队列」真正的业务处理丢给独立的业务线程池。2.2 AsynchronousChannelGroup 的线程分配到底该配多少个 IO 线程AsynchronousChannelGroup 是 AIO 的线程池核心你通过 AsynchronousServerSocketChannel.open(group) 把 broker 的监听套接字挂到 group 上或者通过 AsynchronousSocketChannel.open(group) 把客户端套接字也挂上去。这个 group 里的线程数量直接决定你的并发上限但不是越大越好。我的经验是IO 线程数 CPU 核数 × 2。比如 16 核机器配 32 个 IO 线程每个线程可以同时管理成千上万个 Channel 的读写操作。因为 AIO 的线程不参与数据复制数据搬运是内核完成的线程只负责「被通知」和「发起下一次 IO 操作」。如果你把 IO 线程配到 200 个线程切换的开销反而会吃掉 AIO 带来的收益而且每个线程都有自己的堆栈和本地缓存内存开销也上来了。创建 group 的代码看起来很简单但有一个坑group 的线程工厂一定要设置 Daemon 线程否则你的客户端组件在嵌入到别人的应用时退出主线程后 JVM 不会退出因为 IO 线程还活着AsynchronousChannelGroup group AsynchronousChannelGroup.withFixedThreadPool( Runtime.getRuntime().availableProcessors() * 2, r - { Thread t new Thread(r, mqtt-io-thread); t.setDaemon(true); return t; }); AsynchronousServerSocketChannel server AsynchronousServerSocketChannel.open(group);这段代码里 withFixedThreadPool 的第一个参数是线程数第二个参数是线程工厂。注意我把线程名也定了这在排查问题的时候非常有用——你可以通过 jstack 看到某个 Channel 的读写请求具体卡在哪个 IO 线程上。线程名建议和连接方向、组件类型挂钩比如 broker 端就叫 mqtt-broker-io-thread客户端就叫 mqtt-client-io-thread别偷懒用默认的。2.3 连接生命周期里每一次 IO 事件对应的回调时机MQTT 连接从建立到关闭需要经历 accept、read、write、timeout、close 这几个事件。AIO 下这些事件全是异步回调但你要记住一个原则每次读操作只读一次且必须重新发起下一次读。很多新手在 completed() 里解析完报文就完事了不再调用 read()导致连接静默挂死客户端发什么 broker 都收不到。public class MqttConnection { private final AsynchronousSocketChannel channel; private final ByteBuffer readBuffer ByteBuffer.allocate(8192); public void startRead() { channel.read(readBuffer, this, new ReadCompletionHandler()); } private class ReadCompletionHandler implements CompletionHandlerInteger, MqttConnection { Override public void completed(Integer bytesRead, MqttConnection attachment) { if (bytesRead 0) { close(); // 对端关闭了连接 return; } readBuffer.flip(); // 这里是拆包逻辑从 readBuffer 里解析出完整的 MQTT 报文 MqttPacket packet MqttPacketDecoder.decode(readBuffer); if (packet ! null) { // 丢给业务线程池处理不能在这里直接处理 businessExecutor.submit(() - handlePacket(packet)); } readBuffer.clear(); attachment.startRead(); // 关键必须重新发起下一次读 } Override public void failed(Throwable exc, MqttConnection attachment) { close(); } } }这段代码里最重要的是 startRead() 的自循环调用。completed() 执行完不代表这个连接的数据读完了TCP 是流式协议你可能只读到半个报文所以必须立即再次发起 read让内核继续往 readBuffer 里填数据。还有一个细节readBuffer 的分配。8KB 是一个比较通用的初始值因为 MQTT 的报文体比如订阅的 topic 列表、发布的消息 payload一般不会太大8KB 能覆盖大多数场景。但如果你要传大文件或者批量消息单个报文可能超过 8KB这时候要么在拆包逻辑里做缓冲区扩容要么在解码器里支持分片累积。建议做成分片累积模式每次 read 的数据先放到一个累积缓冲区解码器尝试从累积缓冲区解析出一个完整报文没有就继续等下一次 read。这个模式我在下面的 broker 实现里会展开。3. 搭建百万级 MQTT Broker从 accept 到心跳回收的核心链路3.1 用 AsynchronousServerSocketChannel 写 broker 的 accept 循环Broker 的入口是 accept 操作。和 BIO 的 accept() 不同AIO 的 accept() 也是异步的你在 CompletionHandler 的 completed() 里拿到新的 AsynchronousSocketChannel然后立刻发起下一次 accept同时为新连接初始化读写上下文。有一个容易踩的坑completed() 的 attachment 参数。很多人忽略它直接在回调里用外部变量但外部变量在多连接场景下就是共享状态会串。正确做法是把每个新连接封装成一个连接对象作为下一次 accept 的 attachment 传进去这样每个回调里的数据都是连接私有的。public class MqttBroker { private final AsynchronousServerSocketChannel serverChannel; private final ConnectionManager connectionManager; public void start() throws IOException { serverChannel.accept(null, new AcceptCompletionHandler()); } private class AcceptCompletionHandler implements CompletionHandlerAsynchronousSocketChannel, Void { Override public void completed(AsynchronousSocketChannel clientChannel, Void attachment) { // 先发起下一次 accept否则新连接进不来 serverChannel.accept(null, this); // 为每个连接创建独立的上下文 MqttConnection conn new MqttConnection(clientChannel); connectionManager.register(conn); conn.startRead(); } Override public void failed(Throwable exc, Void attachment) { // accept 失败不能停掉整个服务记录日志后继续 serverChannel.accept(null, this); } } }注意 accept 的 completed() 里我先调用了 serverChannel.accept(null, this)再初始化连接。这个顺序很重要如果你先做连接初始化再做 accept在初始化耗时较长时内核里排队的连接请求就会堆积新连接建立速度会明显下降。正确姿势永远是「先补位再处理当前」。failed() 里同样要重新发起 accept因为临时性的文件描述符耗尽或连接重置是常态不能因为一次失败就退出。3.2 连接状态机从 CONNECT 到 SUBSCRIBE 再到 PUBLISH 的流转每一个 MqttConnection 内部都维护着一个状态机状态机的流转就是 MQTT 协议的核心。协议状态包括等待 CONNECT、已连接待订阅、已订阅可收发消息、正在断线重连。你在 AIO 的读回调里解析出报文后根据当前状态决定执行什么动作。public class MqttConnection { private volatile MqttState state MqttState.IDLE; public void handlePacket(MqttPacket packet) { switch (state) { case IDLE: // 第一个报文必须是 CONNECT if (packet.type() MqttType.CONNECT) { handleConnect(packet); state MqttState.CONNECTED; } else { close(); // 协议违规直接断开 } break; case CONNECTED: switch (packet.type()) { case SUBSCRIBE - handleSubscribe(packet); case PUBLISH - handlePublish(packet); case PINGREQ - handlePing(); case DISCONNECT - close(); } break; } } }这段代码用 Java 17 的 switch 箭头语法如果你还在用 Java 8改回冒号写法就行。需要补充的是客户端的状态流转是「CONNECT - 发送 CONNACK - 等 SUBSCRIBE - 收 PUBLISH」而 broker 的状态流转是「等 CONNECT - 发 CONNACK - 等 SUBSCRIBE - 存订阅关系 - 转发 PUBLISH」。两者是镜像关系你只要把一份状态机代码写对客户端和 broker 都能复用差别只在于一个主动发起、一个被动响应。心跳超时回收是 broker 里最容易出问题的点。每个连接在 CONNECT 报文里会带一个 keepalive 字段单位是秒broker 必须在超过 1.5 倍 keepalive 时间没收到任何报文时断开连接。你不能为每个连接开一个 Timer 线程百万连接就百万个定时器这不现实。要在 broker 里做一个全局的定时扫描器定时遍历所有连接的 lastReadTime把超时的清理掉。3.3 百万连接的存储和检索用哪个容器装连接是个大问题连接状态存在 ConcurrentHashMap 是最直观的做法但百万级容量下有隐患ConcurrentHashMap 的 key 如果是连接 ID 字符串内存里还要存一份 key 的副本百万个 key 就是几十 MB。更关键的是 map 的扩容和并发写入竞争。我一般推荐用 Netty 的 io.netty.util.collection.LongObjectHashMap 或者自研一个基于 Long 到连接对象的映射因为你可以用「连接自增 ID」作为 key避免字符串对象开销。还有一个必须处理的分片问题单台机器百万连接意味着你要对连接做分桶管理。常见做法是每个 IO 线程绑定一个连接分片分片内用数组或链表存连接引用。这样当定时扫描线程去检查心跳时是按分片并行扫描的不会出现一个线程遍历百万对象导致扫描周期过长。分片大小一般取 1024 或 2048数组预分配避免频繁扩容。public class ConnectionManager { private static final int SEGMENT_BITS 10; // 2^10 1024 个分片 private static final int SEGMENT_SIZE 2048; private final MqttConnection[] segments new MqttConnection[1 SEGMENT_BITS]; private final AtomicLong nextConnectionId new AtomicLong(0); public long register(MqttConnection conn) { long id nextConnectionId.incrementAndGet(); int segmentIndex (int) (id ((1 SEGMENT_BITS) - 1)); int slotIndex (int) (id SEGMENT_BITS); // 数组槽位 连接 ID 关联避免全表扫描 segments[segmentIndex * SEGMENT_SIZE slotIndex] conn; conn.setId(id); return id; } }这里用 id 的低 10 位定位分片高位移位后定位分片内的槽位相当于一个二维数组。注意这种固定大小的数组在连接数超过百万后slotIndex 会超出 SEGMENT_SIZE 导致数组越界所以要么启动时根据预估容量计算分片参数要么在 slotIndex 越界时做动态扩展。我对生产环境的建议是预估好你的最大连接数一次性分配够别做动态扩展动态扩展在内存分配时会造成长时间停顿。3.4 心跳扫描和无效连接回收别让僵尸连接占着文件描述符MQTT 的僵尸连接是 broker 的头号杀手。客户端断网后 TCP 层可能半天都发现不了如果你不做应用层心跳检测那些死连接会一直占着文件描述符和内存里的连接对象。百万级场景下每多一个僵尸连接系统可用的文件描述符就少一个直到 reach 上限后新连接全部拒绝。我做的方案是一个后台线程每 5 秒扫描一次全部连接分片。扫描逻辑不是遍历每个连接做系统调用那样太慢而是检查连接对象的 lastReadTime 字段volatile long当前时间 - lastReadTime 超过 keepalive 的 1.5 倍就把连接标记为待回收然后由 IO 线程或专门的清理线程执行 close()。这里有个容易踩坑的细节lastReadTime 的更新位置。必须在 AIO 的 read 回调的 completed() 里第一时间更新不能在业务线程池里更新否则报文已经在 TCP 缓冲区里等了好一会儿才被业务线程更新时间戳不准会把存活连接误判为超时。public void heartbeatScan(long now) { while (true) { long deadline now - keepaliveMillis * 3 / 2; for (MqttConnection conn : allConnections()) { if (conn.lastReadTime deadline) { conn.close(); connectionManager.remove(conn.id()); } } try { Thread.sleep(5000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); return; } } }这段心跳扫描看着简单但 allConnections() 的遍历代价是 O(N)百万连接一次要遍历百万个对象。我用分片后每个分片是独立数组扫描线程可以并行分片扫描。还有一个优化不在扫描线程里直接 close()而是把待回收连接放到一个ConcurrentLinkedQueue里由专门的重启线程去 close。这样避免扫描线程在 close() 时阻塞因为 close() 触发 AIO 的读回调会抛 ClosedChannelException这个异常处理也是要耗时。4. MQTT 客户端组件的正确写法用 AIO 封装发布订阅的异步语义4.1 客户端的 AIO 读循环和 QoS 语义如何映射到回调客户端组件的核心难点不是连接和收发报文而是把 MQTT 的 QoS 语义映射到 AIO 的异步回调上。QoS 0 是即发即弃最简单QoS 1 要求 PUBACK 确认收到确认前消息要留在发送队列里超时重发QoS 2 是四步握手PUBLISH - PUBREC - PUBREL - PUBCOMP状态更复杂。我见过很多人把 QoS 逻辑写在 AIO 的 completed() 回调里这是错的。completed() 里只能做「收到 PUBACK把发送队列里对应的消息标记为已确认」而「判断超时重发」不能在这里做。因为 completed() 的触发时机是 IO 线程你没法在这个线程里 sleep 等超时。public class MqttClientAio implements MqttClient { private AsynchronousSocketChannel channel; private final MqttSession session new MqttSession(); // 保存会话状态和待确认消息 public void publish(String topic, byte[] payload, int qos) { MqttPacket packet MqttPacketFactory.publish(topic, payload, qos); // 写入发送队列后立即触发 channel.write不等业务处理 session.enqueueOutgoing(packet); channel.write(ByteBuffer.wrap(packet.toBytes()), null, new WriteCompletionHandler()); } private class WriteCompletionHandler implements CompletionHandlerInteger, Void { Override public void completed(Integer result, Void attachment) { // 写完成不代表 PUBACK 到了这里不能移除发送队列 MockMqttLog.debug(write completed: {} bytes, result); } Override public void failed(Throwable exc, Void attachment) { session.markChannelBroken(); reconnectOrNotify(); } } }这段代码里的session.enqueueOutgoing和channel.write是两回事。enqueue 是把消息复制到会话的发送队列一般用 ConcurrentLinkedDequewrite 是把字节交给内核。真正让消息从队列里移除的动作发生在「收到 PUBACK」时由读回调解析出 PUBACK 报文后调用session.ack(packet.packetId())完成。QoS 2 的状态机比 QoS 1 麻烦得多。你需要为每个 packetId 维护一个状态发送 PUBLISH 等待 PUBREC - 收到 PUBREC 发 PUBREL 等待 PUBCOMP - 收到 PUBCOMP 结束。这个状态机也必须放在 MqttSession 里不能放在 IO 回调里。4.2 断线重连和会话恢复如何利用 MQTT 的 CleanSession 标志客户端组件的工程难点是断线重连。TCP 断开的发现是滞后的AIO 的回调里可能过了很久才触发 failed()。在你发现连接断开之前用户代码可能已经调用了多次 publish这些消息不能丢尤其 QoS 1/2也不能无限堆积。MQTT 协议的 cleanSession 标志决定了重连后的行为。cleanSession1 时重连后服务端丢弃所有旧会话状态客户端要把未确认的消息重新发送cleanSession0 时服务端保存会话状态重连后可以续传。客户端组件要支持这两个模式并且要在重连时根据标志决定是重建会话还是恢复会话。public void reconnect() { if (session.getCleanSession()) { session.cleanUp(); // 清理所有待确认消息重新走一遍 CONNECT 流程 } else { session.markReconnecting(); // 保留 pendingAck 队列等重连后重新发送 } connectAsync(); } private void connectAsync() { AsynchronousSocketChannel channel createChannel(); channel.connect(new InetSocketAddress(host, port), null, new ConnectCompletionHandler()); } private class ConnectCompletionHandler implements CompletionHandlerVoid, Void { Override public void completed(Void result, Void attachment) { MqttPacket connectPacket MqttPacketFactory.connect( clientId, session.getCleanSession(), keepaliveSeconds); channel.write(ByteBuffer.wrap(connectPacket.toBytes()), null, new ConnectWriteHandler()); } Override public void failed(Throwable exc, Void attachment) { // 指数退避重连避免对 broker 造成重连风暴 scheduleReconnect(); } }重连的指数退避是必修课。如果百万客户端同时断线同时重连broker 瞬间会被 CONNECT 报文打爆。退避策略一般从 1 秒开始每次翻倍最大到 64 秒加上随机抖动±20%避免统一节奏。4.3 背压机制和发送队列上限别让百万客户端的发布压垮 broker客户端组件还有一个隐形坑背压。用户代码调 publish() 是同步的但 AIO 的 write() 是异步的如果业务线程疯狂 publish而内核发送缓冲区已满channel.write() 会一直积压在 IO 线程队列里内存被发送队列撑爆。我一般会在 MqttSession 里设置发送队列的容量上限。比如 QoS 1 的待确认队列最大 1024 条超过后 publish() 返回 false 或者抛异常让业务侧感知到压力。这个对百万级场景尤其重要broker 处理不过来时积累在客户端的内存里比丢在网络里更可怕。public boolean publish(String topic, byte[] payload, int qos) { if (session.outgoingQueue.size() maxOutgoingQueueSize) { // 背压触发拒绝新消息进入 return false; } MqttPacket packet MqttPacketFactory.publish(topic, payload, qos); session.enqueueOutgoing(packet); channel.write(ByteBuffer.wrap(packet.toBytes()), null, new WriteCompletionHandler()); return true; }那 broker 端的背压怎么做broker 接收到客户端的大量 PUBLISH 后要转发给其他订阅者。如果某个订阅者下游客户端消费速度跟不上broker 不能无限往它的 TCP 缓冲区塞数据。我看到很多 broker 的实现是直接往 channel.write() 丢SocketChannel 缓冲区满了就抛异常然后断开连接。这样太粗暴。常规做法是给每个连接单独维护一个发送队列或者叫 pending-write queue如果队列超过水位线比如 8192 条就把这条连接标记为「慢消费者」并考虑断开它或者降级丢弃 QoS 0 消息。这个水位线要根据你的内存预算动态调。5. 避坑指南百万级 MQTT 场景的五个实战踩坑记录5.1 现象连接数到 3 万后 accept 开始超时CPU 没满但吞吐下降原因AsynchronousChannelGroup 的 IO 线程池被读回调里的业务逻辑阻塞了。我的业务线程池里一旦有慢查询比如订阅关系存储内存中的 ConcurrentHashMap 扩容没有注意把耗时操作移出 IO 回调导致 IO 线程排队处理读事件accept 事件排在后面新连接的握手就慢了。解决把 IO 回调里的业务处理全部丢到独立的业务线程池IO 线程只做「拆包 状态机跳转 再发起读」。我后来用了一个 RingBuffer 做「IO 线程到业务线程」的无锁队列把 IO 线程的压力彻底降下来。5.2 现象客户端发送窗口关闭后服务端写回调永远不触发原因TCP 的滑动窗口满了。当你往对端写数据时对端应用层不消费比如客户端在处理业务没调用 read协议栈的发送缓冲区持续堆积直到 window 为 0。此时 channel.write() 的异步操作不会完成你的 WriteCompletionHandler.completed() 一直不执行后续消息全卡在 pening 队列里。解决给写操作加超时。AIO 的 write() 没有内置超时参数我一般启动一个 scheduled task定期检查 pendingWrite 队列的大小和最近一次 write 完成时间。如果超过 30 秒可配没有完成直接关闭连接。那么客户端会感知到「连接异常」自动重连避免永久卡死。5.3 现象内存飙升到 GC 完全失控Full GC 每秒一次原因ByteBuffer 分配过多没有池化。每收到一个报文就 allocate 一个新的 ByteBuffer百万连接在吞吐高峰时瞬间产生大量堆外内存对象。堆外内存是无法被 JVM GC 管理的只有堆内 Buffer 能回收但我的业务把 payload 又复制了一份堆内数组导致堆被撑爆。解决用池化的 ByteBuffer。每次 read 用同一个 Buffer只在下一次 read 前小心清理标记。数据从 Buffer 中解析出 MQTT 报文后payload 直接引用 Buffer 的切片slice不复制。如果业务线程需要长时间持有 payload再复制进堆内数组。这样堆外的对象数大幅减少堆内也因为复用降下来了。5.4 现象broker 重启后客户端一直重连不上服务端报 too many open files原因操作系统的文件描述符限制。百万连接需要百万个文件描述符Linux 默认 ulimit -n 是 1024根本不够用。还有 ephemeral port 范围限制如果 broker 主动发起到下游的连接百万连接会把本机端口耗尽。解决改内核参数。ulimit -n 调高到百万以上net.ipv4.ip_local_port_range调宽net.core.somaxconn调大对于客户端组件要避免主动发起到 broker 的连接使用默认随机端口可以配置复用同一组端口。AIO 本身不是这里的问题但工程上线时一定要验证你的操作系统配额。5.5 现象用 jstack 看不到业务线程但 broker 卡顿严重延时飙高原因这个坑很隐蔽——我在业务线程池里用了同步的 JDBC 操作而业务线程池被数据库连接池耗尽了。当订阅关系或 ACK 处理需要查库时连接池只有 10 个连接而来了 10000 个报文业务线程全部阻塞在获取数据库连接上。IO 线程不断往业务队列丢任务队列积压。解决把存储从 DB 改成内存中的数据结构如 ConcurrentHashMap 存储订阅关系、用 LongObjectHashMap 存储会话并把所有 DB 操作异步化或者批量落盘。MQTT broker 本质上是个内存型系统不要在热路径上做同步 DB 访问。对于持久化要求可以用 WALWrite-Ahead Log的方式异步批量刷新别拿同步 JDBC 硬扛。6. 进阶验证百万级能力的必要步骤与两个核心参数调优6.1 分布式部署中的长连接负载均衡客户端调度策略单台机器扛百万连接很难即使扛住了broker 的单点故障风险也无法接受。在实际落地上我一般会把 broker 做成集群每台 broker 维护一部分长连接客户端连接时根据负载算法分配到不同的 broker。那用 AIO 写的 broker 集群怎么协调核心是「连接粘连」和「主题分区」。连接粘连是指客户端重连时尽量回到同一台 broker这样会话恢复的开销最小主题分区是指某个 topic 的消息固定由某台 broker 或者某个分片来负责其他 broker 收到订阅后要转发给负责该分片的那台。这层逻辑可以用 ZooKeeper 或 etcd 做元数据协调也可以用 Redis Pub/Sub 做消息中转。但要注意的是一旦引入集群组件AIO 带来低延迟收益可能会被跨节点的网络开销抵消一部分所以第一版最好是单机优化到位再加集群。6.2 压测验证自己写一个百万连接的压测脚本你没法真的找一百万个设备但可以用「每客户端多连接」的方式模拟。我给你一个思路在性能测试机上创建 1000 个线程每个线程维护 1000 个 AsynchronousSocketChannel 连接就能模拟 100 万条客户端连接。这里 AIO 的优势是单个客户端进程可以创建大量连接。netstat -an | grep :1883 | wc -l这个命令看当前 broker 上的连接数。但要注意模拟连接并不等于真实业务流量。你要观察三个指标消息吞吐每秒转发消息数、P99 延迟从客户端发布到 broker 转发到订阅者的时间、内存与 GC 状况。压测时要同时跑发布端和订阅端订阅端要消费掉消息否则消息堆积会触发背压逻辑导致压测结果失真。6.3 两个关键调优参数SO_RCVBUF 和 keepalive 扫描粒度第一是 SocketChannel 的接收缓冲区大小。AIO 的内核缓冲区如果太小高吞吐场景下会频繁触发 read 回调CPU 空转如果太大慢消费者会把内存耗尽。我一般设置在 64KB 到 256KB 之间并开启 TCP_NODELAY禁用 Nagle 算法。因为 MQTT 消息通常是小包Nagle 会导致小包被合并延迟升高。clientChannel.setOption(StandardSocketOptions.SO_RCVBUF, 64 * 1024); clientChannel.setOption(StandardSocketOptions.SO_SNDBUF, 64 * 1024); clientChannel.setOption(StandardSocketOptions.TCP_NODELAY, true);第二是心跳扫描线程的粒度。扫描太频繁1 秒浪费 CPU太稀疏30 秒会导致僵尸连接存活过久。我的经验是 5 秒扫描一次乘以 1.5 的 keepalive 容忍系数。如果 keepalive 是 60 秒一个死连接最长 90 秒才被回收这是 MQTT 规范允许的。还有一个验证小技巧上线前用ss -s观察系统 socket 总数用vmstat观察上下文切换用jstat -gcutil盯 Full GC 次数。百万连接场景下Full GC 超过每 10 分钟一次就必须得优化内存结构别再调堆大小了多半是对象复用没做好。这套方案跑下来我的直观体会是Java AIO 的异步回调模型高度契合 MQTT 长连接这种「高并发低峰值」的场景它帮你把线程数压到极低但代价是编程模型从「顺序流」变形为「回调链」。你必须在设计阶段就把以下三件事想清楚——拆包逻辑无状态化、连接对象池化持久化、IO 线程绝不阻塞。这三件事想不透百万级就是压测横幅上的数字不是生产环境里的真实吞吐。我做过的项目里把这三件事做实的 broker 在 16 核机器、128G 内存的配置下稳定扛住了接近 80 万条连接P99 延迟不到 20ms。压在最后的是你的耐心和排查手段——连接数的上升过程里每个量级都有它自己的坑。希望帮到你。本文还有配套的精品资源点击获取