
简介这套基于 Java AIO 实现的 MQTT 客户端与 Broker 组件面向物联网、边缘计算、消息服务器等需要高并发低延迟通信的技术场景适合后端工程师、物联网平台开发者阅读和二次改造。代码完整支持 MQTT v3.1、v3.1.1 与 v5.0同时集成 WebSocket 子协议可兼容 mqtt.js 等前端客户端服务端还提供 HTTP REST API 便于管理并支持遗嘱消息、保留消息等 MQTT 高级特性。集群方面通过 Redis 发布订阅实现消息转发也可自定义消息处理来扩展工程中提供 Spring Boot 快速接入示例支持 GraalVM 编译为本机程序并预留 Prometheus 与 Grafana 监控对接方式整体设计方案比较完整。资源包共 282 个文件以 221 个 Java 源文件为主体配合 Markdown 文档用于说明XML/YAML 文件用于配置另有 HTTP 调试脚本、属性文件等辅助内容压缩包仅 502KB便于快速下载和源码审阅。目前已有 279 人学习/下载从中可以拿到客户端与 Broker 的核心实现、集群方案与监控集成思路是一份高参考价值的 MQTT 实践代码库。1. 百万级 MQTT 连接到底难在哪Java AIO 不是银弹但值得做一条 2C4G 的云主机想撑住 100 万个 MQTT 长连接你最先遇到的可能不是 CPU而是文件描述符和内存。Java AIONIO.2用回调代替阻塞轮询让每个连接无需独占线程这是它敢做百万级 client / broker 的底气但换个角度说AIO 只是把 I/O 搬运交给内核真正决定能不能扛住的是回调线程、ByteBuffer 复用和背压策略。这篇文章按我自己的落地路径写先看线程模型再写 client 和 broker最后给压测与避坑清单。适合正准备用 Java 写 MQTT 接入层、或者想把 Netty 方案换掉的人。2. 选型先看线程模型Java AIO、NIO 与 Netty 在百万连接下的真实差异2.1 Java AIO 的回调模型谁在替你执行 I/O 回调Java AIO 在 JDK 1.7 进入 NIO.2底层在 Linux 上通常走 epoll在 Windows 上走 IOCP。与 NIO 的Selector不同AIO 把“可读 / 可写”的通知包装成CompletionHandler回调内核把数据拷进ByteBuffer之后才发起回调。也就是说应用不用自己维护OP_READ就绪状态连接也不需要一个线程死等。生产环境里的第一件事是显式创建AsynchronousChannelGroup否则 JVM 会用自己的默认线程池大小不可控而且和业务线程混在一起很难排查。AsynchronousChannelGroup group AsynchronousChannelGroup .withFixedThreadPool( Runtime.getRuntime().availableProcessors() * 2, r - new Thread(r, aio-callback)); AsynchronousSocketChannel channel AsynchronousSocketChannel.open(group);这段代码的逻辑很直白所有 AIO 回调都会丢进这个固定线程池。核心参数是availableProcessors() * 2。为什么不是 1:1因为回调里往往还要做少量协议解析线程多一点能降低抖动但也绝对不能开成几百个线程回调任务大多是短命任务太多线程反而增加上下文切换。回调模型是双刃剑回调里如果出现阻塞比如同步写 DB、打日志、加锁等线程池会被瞬间打满。特别要注意的是AIO 回调的异常如果没被CompletionHandler.failed接住异常会进入线程池的兜底逻辑连接挂了你都不知道。我一般会为回调线程池单独声明一个未捕获异常处理器把所有异常打到一个独立日志文件里。2.2 MQTT 报文边界与半包为什么解码必须先做长度字段MQTT 的报文不是按行分割的也没有 HTTP 那种空行。它靠固定报头里的“剩余长度”字段切包。剩余长度是变长编码每个字节低 7 位是有效值最高位是“是否还有下一个字节”最多占 4 字节。所以收到 TCP 字节流后第一步永远是解析剩余长度第二步是攒够这么多字节再消费一个完整报文。这里的核心难点在“等”。TCP 是流协议一次read回调可能只到了半个报文也可能一次到了好几个报文。如果拿到半包就丢给上层MQTT 路由就全乱了。正确的做法是读回调拿到数据先flip进入读模式尝试从当前位置解析固定报头和剩余长度判断缓冲区剩余字节是否足够承载整个报文不够就回退到起始位置compact等待下一次read回调。这个逻辑会在第 3 章给出完整代码。这里想强调一点剩余长度本身也可能跨多个 TCP 分段比如剩余长度编码占了 3 个字节TCP 只先到了 2 个字节这时同样要回退不能直接报协议错误。2.3 项目里选 AIO 还是 Netty一张对比表说清楚我在决定用 Java AIO 之前其实先写了一版 Netty 的原型。Netty 很成熟但团队当时想做一个零依赖的长连接接入组件并且连接上的业务逻辑很简单解码 MQTT、查订阅、转发。这个场景用 Netty 会引入大量用不到的能力所以我最终选择了 AIO。对比项Java AIONetty线程模型内核写完/读完才回调固定回调线程池Reactor 多线程EventLoop 与连接绑定编解码手写半包/粘包组装自带LengthFieldBasedFrameDecoder等内存池无内建池需要自己复用ByteBufferPooledByteBuf成熟依赖体积JDK 原生零额外 jar 包至少引入 netty-transport 等模块排查成本回调异常容易吞需要自己打日志社区资料多异常栈相对清晰适合场景连接接入层、协议简单、想控依赖复杂网关、多协议、代理、TLS 场景这张表的结论是如果只是做 MQTT 接入和转发Java AIO 完全够用如果要同时支持 WebSocket、SCTP、自研协议扩展直接上 Netty 更划算。标题既然限定在 Java AIO后面就全部围绕 AIO 展开。3. 用 Java AIO 写一个能扛压的 MQTT client 组件连接、订阅、重连三板斧3.1 最小可运行的 AIO Client建立连接与 CONNECT 报文先看最基础的连接建立。连接是异步的connect方法立即返回真正完成连接在CompletionHandler.completed里。这个回调里要马上发 CONNECT 报文然后开始第一次read。public class AioMqttClient { private AsynchronousSocketChannel channel; private final ByteBuffer readBuffer ByteBuffer.allocate(4096); private volatile boolean connected false; public void connect(String host, int port, String clientId, int keepAlive) throws Exception { AsynchronousChannelGroup group AsynchronousChannelGroup .withFixedThreadPool(4, r - new Thread(r, mqtt-client-aio)); channel AsynchronousSocketChannel.open(group); channel.connect(new InetSocketAddress(host, port), null, new CompletionHandlerVoid, Void() { Override public void completed(Void result, Void attachment) { connected true; sendConnectPacket(clientId, keepAlive); channel.read(readBuffer, null, new ReadCompletionHandler()); } Override public void failed(Throwable exc, Void attachment) { scheduleReconnect(host, port, clientId, keepAlive); } }); } private void sendConnectPacket(String clientId, int keepAlive) { ByteBuffer buffer ByteBuffer.allocate(256); // CONNECT 报文的固定报头控制包类型为 0x10 buffer.put((byte) 0x10); int remainingStart buffer.position(); buffer.put((byte) 0x00); // 剩余长度占位后面回填 buffer.putShort((short) 4); buffer.put(MQTT.getBytes()); buffer.put((byte) 4); // 连接标志0x02 表示 clean session 为 1 buffer.put((byte) 0x02); buffer.putShort((short) keepAlive); buffer.putShort((short) clientId.length()); buffer.put(clientId.getBytes()); int remaining buffer.position() - remainingStart - 1; buffer.put(remainingStart, (byte) remaining); buffer.flip(); channel.write(buffer, null, new WriteCompletionHandler()); } }这里有两个参数一定要解释清楚。keepAlive的单位是秒客户端要在keepAlive / 2的间隔发送 PINGREQ否则 broker 会在 1.5 倍 keepAlive 时间内等不到报文时杀掉连接。clean session 0x02的位是连接标志里的 bit1置 1 表示服务端不保留 session重连后订阅关系全部作废这样做重连逻辑最简单。注意上面的剩余长度回填只适用于一个字节的剩余长度实际生产里要用变长编码去写否则报文超过 127 字节就错了。这个细节放到第 5 章避坑里讲因为真的有人在这里翻车。3.2 读写回调与 MQTT 解码器的拼装读回调是客户端最核心的部分。每次read完成completed里传进来的Integer是本次读取的字节数小于 0 表示对端关闭了连接。拿到数据后先decode再发起下一次read这是 AIO 长连接的标准姿势。private final class ReadCompletionHandler implements CompletionHandlerInteger, Void { Override public void completed(Integer readBytes, Void attachment) { if (readBytes 0) { close(); return; } decode(readBuffer); channel.read(readBuffer, null, this); } Override public void failed(Throwable exc, Void attachment) { close(); } } private void decode(ByteBuffer buffer) { buffer.flip(); while (buffer.remaining() 0) { int start buffer.position(); if (buffer.remaining() 2) { break; } int first buffer.get() 0xFF; int multiplier 1; int remainingLength 0; int current; int headerLen 1; do { current buffer.get() 0xFF; remainingLength (current 0x7F) * multiplier; multiplier 7; headerLen; } while ((current 0x80) ! 0); int totalLen headerLen remainingLength; if (buffer.remaining() remainingLength) { buffer.position(start); break; } buffer.position(start); byte[] packet new byte[totalLen]; buffer.get(packet); dispatch(packet); } buffer.compact(); }这段解包代码是 MQTT 客户端见高下的地方。三个参数要重点说readBuffer.allocate(4096)表示单次读取的最大字节数。MQTT 报文理论上可以很大比如 256MB 的 PUBLISH但客户端内存有限我一般限制单条报文最大不超过readBuffer.capacity()超过直接关闭连接避免恶意报文撑爆内存。remainingLength解析最多 4 个字节当multiplier 1 21时必须踢掉连接这是协议硬性要求。最后一行compact()把未消费的字节移到缓冲区头部下一次 read 继续往后写。如果没有这行粘包数据会丢失。3.3 重连与心跳的坑CONNECT 不能重发同一个报文的 packetId连接断开后不能无限立刻重连需要指数退避。常见做法是 1 秒起步乘以 2封顶 60 秒再加一个 0 到 1 秒的随机抖动防止 10 万个客户端同时在断线后重连造成连接风暴。心跳的调度用独立单线程定时器不能用 AIO 回调线程来做定时否则回调线程被定时任务占住数据读写就没线程了。ScheduledExecutorService scheduler Executors.newSingleThreadScheduledExecutor(); scheduler.scheduleAtFixedRate(() - { if (!connected) { return; } long elapsed System.currentTimeMillis() - lastResponseTime; if (elapsed keepAlive * 1000L * 1.5) { close(); scheduleReconnect(host, port, clientId, keepAlive); return; } sendPingReq(); }, keepAlive * 1000L / 2, keepAlive * 1000L / 2, TimeUnit.MILLISECONDS);这个逻辑里有两个坑。第一lastResponseTime必须在收到 CONNACK 和 PINGRESP 时更新而不是收到任意报文就更新否则 broker 已经失联但客户端还傻等。第二重连时最好不要复用同一个clientId去连续 CONNECT尤其是clean session 0的场景。如果 broker 上还留着旧连接的 sessiion新连接会触发旧连接断开这个动作在 broker 里是同步踢连接可能阻塞你重连的线程。我习惯在重连时给clientId加一个递增后缀直到 CONNACK 成功后再恢复原样。这样做的好处是服务端 session 状态可控缺点是 clean session 0 时旧 session 会永久残留所以生产环境我的建议是客户端全部走 clean session 1由业务侧补订阅。4. Broker 侧路由基于 AIO 实现百万连接的消息分发与背压4.1 ServerSocketChannel 接收连接与连接注册表Broker 的骨架是AsynchronousServerSocketChannel。accept 回调里要立刻发起下一次 accept否则内核不会再触发回调这是 AIO 服务端最常见的新手坑。public class AioMqttBroker { private AsynchronousServerSocketChannel server; private final ConcurrentHashMapString, ClientContext clients new ConcurrentHashMap(); public void listen(int port) throws IOException { AsynchronousChannelGroup group AsynchronousChannelGroup .withFixedThreadPool(Runtime.getRuntime().availableProcessors() * 2, r - new Thread(r, broker-aio)); server AsynchronousServerSocketChannel.open(group) .bind(new InetSocketAddress(port), 4096); server.accept(null, new CompletionHandlerAsynchronousSocketChannel, Void() { Override public void completed(AsynchronousSocketChannel ch, Void attachment) { server.accept(null, this); handleNewConnection(ch); } Override public void failed(Throwable exc, Void attachment) { // 文件描述符耗尽或端口耗尽时睡眠后继续 accept } }); } }bind的第二个参数是 backlog我一般设 4096这是内核 accept 队列的长度。如果连接风暴瞬间涌入超过 backlog内核会直接丢弃 SYN客户端表现为连接超时。连接注册表是百万连接 broker 的第一个瓶颈。ConcurrentHashMapString, ClientContext是基础但要注意三点clientId冲突时必须主动踢掉旧连接查找、插入、删除都要 O(1)同一时刻可能几百个线程并发操作这个 Map不能用全局锁。我见过有人用Collections.synchronizedMap包一层结果压测时锁竞争直接打满 CPU。4.2 订阅树与消息路由Topic 匹配不要用正则每个 PUBLISH 进来都要匹配所有订阅者的 topic。如果订阅很多正则匹配会在高并下发消息时把 CPU 打爆。正确的做法是维护一棵订阅树按 topic 的/层级展开。class SubscriptionTrie { private final MapString, SubscriptionTrie children new ConcurrentHashMap(); private final MapClientContext, Integer subscribers new ConcurrentHashMap(); public void subscribe(String[] levels, int index, ClientContext ctx, int qos) { if (index levels.length) { subscribers.put(ctx, qos); return; } children.computeIfAbsent(levels[index], k - new SubscriptionTrie()) .subscribe(levels, index 1, ctx, qos); } public void publish(String[] levels, int index, ListClientContext targets) { if (index levels.length) { targets.addAll(subscribers.keySet()); return; } String level levels[index]; SubscriptionTrie plus children.get(); if (plus ! null) { plus.publish(levels, index 1, targets); } SubscriptionTrie exact children.get(level); if (exact ! null) { exact.publish(levels, index 1, targets); } SubscriptionTrie hash children.get(#); if (hash ! null) { collectAll(hash, targets); } } }这段代码是最简订阅树但已经能处理单层通配和#多层通配。注意几个细节#只出现在叶子层匹配剩余所有层$开头的 topic 是系统级不能匹配通配符同一客户端重复订阅时应用新 QoS 覆盖旧值而不是简单加进 Set。发布路由的性能指标是一条消息只遍历路径上的节点而不是遍历全量订阅。100 万个订阅时如果 topic 层级深为 3那么每次发布最多比较 3 层 Map 查找这是可控的。如果换成正则一次发布要扫全量订阅直接没法用。4.3 背压与内存保护写失败时该丢弃还是断开Broker 给客户端下发消息时最怕的不是 CPU而是内存被写缓冲占满。客户端假死或者网络分区TCP 发送缓冲区会先写满这时 AIO 的 write 回调不再触发而发布者还在往这个连接塞消息内存就会无限涨。我常用的方案是给每个连接配一个有界写队列队列满时根据 QoS 做不同处理。public class ClientContext { private final ArrayBlockingQueueByteBuffer writeQueue new ArrayBlockingQueue(1024); private final AtomicBoolean writing new AtomicBoolean(false); public boolean enqueue(ByteBuffer packet) { if (!writeQueue.offer(packet)) { return false; } if (writing.compareAndSet(false, true)) { flush(); } return true; } private void flush() { ByteBuffer packet writeQueue.poll(); if (packet null) { writing.set(false); return; } channel.write(packet, null, new CompletionHandlerInteger, Void() { Override public void completed(Integer result, Void attachment) { if (packet.hasRemaining()) { channel.write(packet, null, this); } else { flush(); } } Override public void failed(Throwable exc, Void attachment) { close(); } }); } }这里的核心是ArrayBlockingQueue(1024)和AtomicBoolean。前者限制单个连接最多积压 1024 个待发送报文后者保证同一时刻只有一个线程在写这个 socket否则并发 write 会导致线程安全问题和报文错乱。队列满时的策略QoS 0 可以直接丢弃这是协议允许的QoS 1 和 QoS 2 不能丢否则破坏“至少一次”语义。我的做法是断开连接让客户端走重连重发同时把该连接标记为高延迟告警出来排查网络。5. 避坑百万连接场景下最常见的 5 个翻车点5.1 读半包与粘包解包没回退现象压测时消息出现错乱有时候一条消息被拆成两条 dispatch有时候两条消息粘成一条客户端直接协议错乱掉线。原因read 回调里拿到多少字节就 dispatch 多少字节没有按 MQTT 剩余长度等完整包。TCP 是流半包和粘包是常态。解决用 3.2 节给到的decode流程先解析剩余长度数据不足时position回退到报文起始位置再compact。我还会在 decode 里加一个最大报文长度校验超过 4096 字节直接断开防恶意大包拖垮内存。5.2 回调线程里做耗时操作线程池被打满现象CPU 没满但连接纷纷超时定期重连日志里出现大量回调线程堆积线程池队列涨到几万。原因在 read 回调里直接同步写了 Redis、MySQL 或打印大量日志。AIO 回调线程池只有 N 个一个回调卡住 200msN 个连接就把线程池占满了。解决回调里只做协议解析和入队把执行业务逻辑的任务丢到独立的ThreadPoolExecutor。这个业务线程池单独设置队列上限和拒绝策略不能把 broker 主流程拖死。5.3 写缓冲无限增长内存 OOM现象broker 堆内存正常但堆外内存持续上涨最后进程 OOM kill。客户端没有断开发布流量也不大。原因某个客户端假死TCP 发送窗口降为 0channel.write的 ByteBuffer 永远发不出去而业务线程还在往该连接塞新报文。解决给每个连接加有界写队列队列满时按 QoS 丢弃或断开。另一道防线是统计连接积压字节数超过阈值比如 16MB直接强制关闭避免单条连接拖垮整个进程。5.4 保活心跳与平台故障超时判定过严现象线上网络抖动客户端在线但 broker 判定离线一断开就是几百上千条连接同时重连形成心跳风暴。原因keepAlive 设 60 秒broker 的判定阈值也设 60 秒抖动只持续 90 秒就误杀了一批。解决服务端超时阈值设为 keepAlive 的 1.5 倍并放到独立的调度线程里扫描不能和 AIO 回调线程混用。如果 JVM GC 偶尔有长停顿阈值还要再放宽一些或者先用 ZGC 把 STW 压下去再上线。5.5 压测客户端先把自己压垮句柄和端口不够现象压测 10 万连接时broker 端还好好的压测机先报Too many open files随后建连速率掉到每秒几十个。原因Linux 默认ulimit -n只有 1024单机可用临时端口只有 32768 到 60999 大约 28000 个达不到百万级 client 压测需求。解决压测机调大文件句柄并放宽本地端口范围。ulimit -n 1048576 sysctl -w net.ipv4.ip_local_port_range1024 65000 sysctl -w net.ipv4.tcp_tw_reuse1如果单机端口仍然不够要么用多台压测机分摊要么让压测客户端复用连接而不是每次重连。这也暴露了一个现实百万级 MQTT client 的瓶颈往往不在 broker而在压测环境本身。6. 从 10 万到 100 万系统参数调优与上线前验证Broker 代码写得再好系统参数不对也扛不住百万连接。先看一组我上线前固定要检查的内核参数参数推荐值作用fs.file-max419430 以上全局文件句柄上限net.core.somaxconn4096accept 队列长度配合 backlognet.ipv4.tcp_mem按内存调整TCP 内存水线避免被 OOM killer 盯上net.ipv4.tcp_tw_reuse1快速复用 TIME_WAIT 端口net.ipv4.ip_local_port_range1024 65000客户端出端口范围压测机必调JVM 侧我给两条具体建议。第一Xms和Xmx设相等连接相关的 Channel 上下文大多是堆内对象避免堆扩容引发长 GC。第二ByteBuffer.allocateDirect使用堆外内存要显式设-XX:MaxDirectMemorySize否则 100 万连接每人一个 4KB 直缓冲就是 4GB默认值会超。上线前验证我习惯做两件事。第一件是断线风暴把 broker 连上 1 万测试客户端然后批量杀掉一半连接观察回调线程池的堆栈和clients表是否正确清理。第二件是连接缓慢建连用自研 MQTT client 组件从 0 到 10 万连接匀速拉起来同时用jstat盯 Young GC 频率和Direct Buffer用量如果 GC 频率明显上升说明连接上下文或者 ByteBuffer 分配有泄漏。我做这类项目积累的最大教训是AIO 回调线程池永远要留有 20% 以上余量。每加一个业务功能就多占一块回调时间时间一长连接就悄悄超时。所以我的习惯是每次发版前都跑一轮连接保活测试确认回调线程池平均利用率不超过 70%并且每个连接的写队列积压数稳定在个位数。这个数据比最后能不能撑住 100 万连接更重要因为它决定了你能不能持续撑住。希望帮到你。本文还有配套的精品资源点击获取