简介:基于 Java 异步 IO(AIO)技术打造的高性能消息队列遥测传输(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 CompletionHandler<Integer, 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 CompletionHandler<AsynchronousSocketChannel, 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 CompletionHandler<Integer, 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 是把消息复制到会话的发送队列(一般用 ConcurrentLinkedDeque),write 是把字节交给内核。真正让消息从队列里移除的动作发生在「收到 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 标志决定了重连后的行为。cleanSession=1 时,重连后服务端丢弃所有旧会话状态,客户端要把未确认的消息重新发送;cleanSession=0 时,服务端保存会话状态,重连后可以续传。客户端组件要支持这两个模式,并且要在重连时根据标志决定:是重建会话还是恢复会话。
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 CompletionHandler<Void, 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 访问。对于持久化要求,可以用 WAL(Write-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。压在最后的是你的耐心和排查手段——连接数的上升过程里,每个量级都有它自己的坑。希望帮到你。
本文还有配套的精品资源,点击获取