简介:压缩包内是一份基于Java与Netty的高并发高可用MQTT消息代理服务端完整工程源码,主要面向需要自建大规模连接的消息平台,或者深入理解消息协议实现的Java服务端工程师。项目使用Netty完成通信层及协议报文解析,使用nutzboot统一管理依赖注入和属性配置,使用Redis支撑消息缓存与集群节点协作,Kafka作为可选的消息代理模块,整体可以轻松支撑十万级并发连接,并已在实际生产环境中长期运行,具备较好的参考与复用价值。压缩包共收录91个文件,其中以59个Java源文件为主,同时包含若干XML、YAML等配置文件,文本说明、证书密钥、系统参数调优说明等辅助材料,整个压缩包仅为231KB,工程内部分为Broker核心、认证、存储、公共模块、客户端等清晰目录,便于按功能阅读和二次开发。目前已有1888人浏览与学习,资料热度与认可度较高,适合希望在Netty长连接服务、消息队列接入、集群缓存同步等方面获取落地经验的开发者。
1. 从“10 万连接”反推 Netty MQTT broker 的模型边界
“10 万并发”这五个字,在 MQTT broker 里比在 Web 网关里更需要拆开看:它通常是 10 万条 TCP 长连接,而不是 10 万 QPS 消息吞吐。一条 MQTT 连接大半时间在等待心跳和平台下发,每秒真正产生上行报文的设备往往只有几千条。所以用 Java + Netty 实现 broker,第一个要立住的模型是“事件驱动 + 少量线程”,而不是“一个连接一个线程”;Netty 管理 10 万个 Channel 不靠堆线程,靠多路复用。真正决定成败的是连接注册、会话状态、QoS 消息存储、集群节点之间的路由扩散,以及生产环境里文件描述符、内存水位、GC 停顿这些容易被低估的参数。这篇内容面向已在 Java 后端实操过 Netty、准备把 MQTT broker 落到生产环境的工程师,也可以拿来回答 netty 面试题里常见的 EventLoop 模型问题。
2. Netty 的线程模型与 MQTT 解码器怎么搭配才不先撞瓶颈
自研 MQTT broker,选 Netty 的理由通常很直接:业务团队是 Java 栈,又要对 MQTT 报文的鉴权、租户隔离、私有属性做深度定制。与其包一层 EMQX 插件,不如在 Netty 的 handler 链里直接控制每个环节。先想清楚:MqttDecoder 负责把 TCP 字节流切成一帧帧 MQTT 报文,MqttEncoder 负责反向编码,真正处理连接状态的是自定义业务 handler。只要这一步设计对,10 万连接在单进程里是可行的。
2.1 用 NIO 管理连接:worker 线程数不等于连接数
我一般会先按下面这个骨架把 broker 跑起来,再逐步加心跳、鉴权和路由。代码里的 bossGroup 只做 accept,workerGroup 负责 channel 的读写,线程数量按 CPU 核数来,不按连接数来。
EventLoopGroup bossGroup = new NioEventLoopGroup(1); EventLoopGroup workerGroup = new NioEventLoopGroup(0); try { ServerBootstrap bootstrap = new ServerBootstrap(); bootstrap.group(bossGroup, workerGroup) .channel(NioServerSocketChannel.class) .option(ChannelOption.SO_BACKLOG, 1024) .option(ChannelOption.SO_REUSEADDR, true) .childOption(ChannelOption.SO_KEEPALIVE, true) .childOption(ChannelOption.TCP_NODELAY, true) .childOption(ChannelOption.WRITE_BUFFER_WATER_MARK, new WriteBufferWaterMark(512 * 1024, 1024 * 1024)) .childHandler(new ChannelInitializer<SocketChannel>() { @Override protected void initChannel(SocketChannel ch) { ch.pipeline().addLast("idle", new IdleStateHandler(90, 0, 0)); ch.pipeline().addLast("decoder", new MqttDecoder(1048576)); ch.pipeline().addLast("encoder", MqttEncoder.INSTANCE); ch.pipeline().addLast("broker", new MqttBrokerHandler()); } }); ChannelFuture future = bootstrap.bind(1883).sync(); future.channel().closeFuture().sync(); } finally { workerGroup.shutdownGracefully(); bossGroup.shutdownGracefully(); }这段代码里,NioEventLoopGroup(1)只留一个线程做 accept,足够应付大量新建连接;NioEventLoopGroup(0)让 Netty 按默认值创建 worker 线程,默认大约为 CPU 核数的两倍。10 万连接下如果每个 EventLoop 都绑定大量 channel,线程不会成为瓶颈,瓶颈往往出在单个 channel 上的阻塞操作。生产环境在 Linux 上应把NioEventLoopGroup换成EpollEventLoopGroup,channel 对应换成EpollServerSocketChannel,能减少一次从 select 到 epoll 的系统调用开销。SO_BACKLOG是 accept 队列长度,只写 1024 还不够,后面要配合操作系统参数一起看。常见参数可按下表设置:
| 参数 | 建议值 | 说明 |
|---|---|---|
SO_BACKLOG | 1024 或 2048 | TCP accept 队列长度,超过net.core.somaxconn时以内核为准 |
SO_REUSEADDR | true | 快速重启 broker,避免端口处于 TIME_WAIT 时无法 bind |
TCP_NODELAY | true | MQTT 报文通常很小,禁用 Nagle 可减少 40ms 延迟 |
SO_KEEPALIVE | true | 先交给内核探测死链,业务层心跳继续自己处理 |
WRITE_BUFFER_WATER_MARK | 512KB / 1MB | 限制单连接在慢客户端场景下的积压缓冲,防止内存被拖垮 |
这里的SO_KEEPALIVE不是替代 MQTT 层的心跳,它只是兜底。真正决定 broker 是否判断设备离线的,是后面的IdleStateHandler。
2.2 Handler 链里的 MqttDecoder、IdleStateHandler 与业务线程池
MqttDecoder()默认限制报文最大字节数,生产环境要按业务放行,否则 1MB 的固件升级报文会被直接断开。上面代码里传入 1048576,表示单帧 MQTT 报文最大 1MB。实际设备上报的大报文场景不多,但下发固件时常见,建议做成配置项。IdleStateHandler(90, 0, 0)表示读空闲 90 秒触发事件;MQTT 的 keepalive 通常由设备侧决定,如果客户端声明 60 秒,broker 最多等 90 秒就该断开,这是 1.5 倍规则的来源。在userEventTriggered里收到IdleStateEvent时关闭连接,是很多 netty 面试题的答法,也是这里唯一正确的处理方式:
@Override public void userEventTriggered(ChannelHandlerContext ctx, Object evt) throws Exception { if (evt instanceof IdleStateEvent) { log.warn("client {} idle timeout, close", ctx.channel().remoteAddress()); ctx.close(); } else { super.userEventTriggered(ctx, evt); } }要注意,这里的超时判断只解决“设备不发送任何数据”的情况。如果设备持续发送错误报文,decode 会走异常分支,需要在exceptionCaught里单独计数,连续多次解析失败同样关闭连接。
业务处理上有一条红线:不要在 Netty 的 EventLoop 线程里做数据库查询、RPC、Redis 同步读写。10 万连接对应大约几十个 worker 线程,任何一个线程因为同步等待卡住几百毫秒,它下面的几千条连接都会出现延迟抖动。常见做法是在 broker handler 里只解析报文、变更内存状态,再交给单独的业务 executor 或消息队列异步处理,处理完再写回对应的 channel。会话 ID 与 channel 的对应关系见第三章。
2.3 慢客户端背压:连接级写水位与 writability 规则
最容易被忽略的问题是慢客户端。设备用 4G 网络信号差,TCP 接收窗口越来越小,如果 broker 持续向这条连接下发消息,Netty 的写出缓冲会一直增长,最后吃掉整块堆外内存。上面的WRITE_BUFFER_WATER_MARK就是第一道闸门:当一条连接未发送的字节数超过高水位 1MB,channel 变成不可写,业务代码需要停止继续投递,等它降到低水位以下再继续。实现上需要检查channel.isWritable(),不要无脑writeAndFlush:
if (ctx.channel().isWritable()) { ctx.writeAndFlush(packet); } else { pendingOfflineQueue.offer(packet); // 等 channelWritabilityChanged 回调后继续发送 }注意这里的pendingOfflineQueue最好不要用无界队列,否则积压过多还是会把 broker 拖垮。更常见的处理是直接落盘,或者走 QoS1 的待确认队列,等客户端恢复后重新投递。这也是 MQTT 场景和普通 TCP 长连接最大的区别:消息不一定要立刻发出,但必须被可靠记录。
3. 连接注册、会话保持和 QoS 消息的存储设计
MQTT 和 HTTP 最大的差异在连接身份:clientId 是设备在业务层的唯一标识。broker 收到 CONNECT 报文后,需要把 clientId 与 ChannelHandlerContext 绑定,后续所有下发都通过这张表找到对应 channel。同时还要考虑同一个 clientId 重复上线、异常断开后残留 channel、session 是否继续保存订阅关系。
3.1 用 clientId 管理连接:重复登录与清理
我用一个 ConcurrentHashMap 保存 clientId 到 ChannelHandlerContext 的映射,同时用另一个 map 保存 session 状态。注册逻辑要处理“旧连接还活着”的情况,不能简单地put覆盖,否则旧连接会继续收到消息,造成消息乱序。
public class ConnectionRegistry { private final ConcurrentHashMap<String, ChannelHandlerContext> channelMap = new ConcurrentHashMap<>(); private final ConcurrentHashMap<String, SessionState> sessionMap = new ConcurrentHashMap<>(); public boolean attach(String clientId, ChannelHandlerContext ctx) { while (true) { ChannelHandlerContext old = channelMap.putIfAbsent(clientId, ctx); if (old == null) { return true; } if (!old.channel().isActive()) { if (channelMap.replace(clientId, old, ctx)) { old.close(); return true; } continue; } return false; // 同一 clientId 已在线,拒绝新连接 } } public void detach(String clientId, ChannelHandlerContext ctx) { channelMap.remove(clientId, ctx); } }这里用remove(key, value)而不是remove(key),是为了防止“客户端 A 断线,客户端 A 的新连接已经绑定成功,旧连接的断开事件又把新映射删掉”这类竞态。生产环境还要在外部存储里放一层设备会话注册信息,用于跨 broker 节点判断重复登录;单机阶段 ConcurrentHashMap 足够。需要留意的是,绝对不要在 EventLoop 线程里遍历整个 channelMap 做广播,否则有 10 万连接时,一次全量遍历会让所有 worker 线程都阻塞在锁竞争上。正确做法是按订阅关系维护 topic 到 clientId 集合的索引,或直接用支持通配符的 TopicTrie。
3.2 QoS 0/1/2 的存储策略与离线消息
MQTT 协议里的消息质量是 broker 必答题。QoS 0 不用持久化,转发失败直接丢;QoS 1 要求 broker 收到 PUBACK 后才能确认完成,因此消息必须先落到“待确认队列”;QoS 2 的流程更长,PUBREC、PUBREL、PUBCOMP 四个报文构成状态机,任何一个状态丢失都会造成重复或丢失。存储策略可以按下表简化:
| QoS | 语义 | Broker 需要做的处理 |
|---|---|---|
| QoS 0 | 最多一次 | 内存中转发,失败不重试,不需要落盘 |
| QoS 1 | 至少一次 | 持久化到待确认队列,收到 PUBACK 后删除,重发时带 DUP |
| QoS 2 | 恰好一次 | 保存报文 PUBLISH 状态机和 packetId,四步握手完成后清除 |
10 万连接场景下,QoS 1 消息是主要压力来源。如果所有消息都写 RocksDB,写放大不可小觑;如果只在内存里放,broker 一重启就丢消息。我通常的做法是:默认 QoS 0 和大部分遥测数据不落盘,只有订阅端不可达且 session 保持时,才把离线消息写入 RocksDB。对在线客户端,先放一个内存队列,收到 PUBACK 后删除;队列长度超过阈值再落盘。这样既能保证消息不丢,也不会让磁盘写成为性能瓶颈。离线消息要有 TTL,否则大量设备长期离线,broker 会把磁盘写满。
3.3 在 CONNECT 包上做鉴权,别把校验漏到业务消息里
总有人问 Netty Websocket 怎么做鉴权,MQTT 的答法类似:鉴权必须发生在连接建立阶段,也就是 CONNECT 到 CONNACK 之间。在 broker 的 handler 里,第一条业务报文必须是 CONNECT,校验不通过就回错误码并关闭连接。常见校验方式是用 MQTT 报文里的 username/password 字段对接内部 token 服务,或使用证书中的 clientId 做设备级校验。
MqttConnectMessage connect = (MqttConnectMessage) msg; String clientId = connect.payload().clientIdentifier(); byte[] passwordBytes = connect.payload().passwordInBytes(); String username = connect.payload().userName(); boolean authorized = accessControl.check(clientId, username, passwordBytes); MqttConnAckMessage ack = MqttMessageBuilders.connAck() .sessionPresent(false) .returnCode(authorized ? MqttConnectReturnCode.CONNECTION_ACCEPTED : MqttConnectReturnCode.CONNECTION_REFUSED_BAD_USER_NAME_OR_PASSWORD) .build(); ctx.writeAndFlush(ack); if (!authorized) { ctx.close(); return; }这里有一个容易踩的坑:不要在channelRead收到第一条消息时直接调用远程 HTTP 认证接口,否则一个慢认证服务会卡住 EventLoop。常见做法是把待认证连接放进“半开连接池”,认证请求异步发出,等回调结果后再回 CONNACK;超时时间一般 5 秒上下,超过即断开。这样即使用户服务抖动,broker 本身不会跟着退化。
4. 高可用必须处理的三件事:会话粘滞、集群广播、持久化
单机 Netty broker 能扛住连接数,不等于能上线。生产环境至少要处理三个问题:网络抖动时连接切换到另一台 broker,客户端会话怎么找到;一条 PUBLISH 消息到达某个节点后,订阅者在别的节点上,怎么把消息送过去;节点宕机后,未确认的 QoS1 消息如何恢复。下面按我常用的套路拆开。
4.1 集群入口:按 clientId 做会话粘滞
常见做法是在 broker 前面加一层四层负载均衡,但 MQTT 长连接不能随便负载均衡。如果连接 A 落在 node1,连接断开后重连落在 node2,而 node1 上的 session 没有同步过去,客户端要么重新订阅,要么丢失离线消息。因此入口层需要按 clientId 做哈希,让同一个设备的连接始终落在同一个 broker 节点上。实现时可以根据 MQTT 报文中解析出的 clientId 计算哈希,也可以用一致性哈希;只按源 IP 哈希不行,因为一个 NAT 后面可能有大量设备,分布会不均衡。只有节点状态同步做得非常好时,才允许自由调度,否则 session 漂移会成为事故源头。
4.2 集群广播:用发布订阅通道做路由扩散
如果 node1 上有客户端发布一条消息,订阅者可能在 node2 上。最简单且最常见的做法是引入一个独立的发布订阅通道,每个 broker 节点订阅同一个 topic;收到本机用户请求后,把编码后的消息广播到通道里,所有节点收到后再判断本机是否存在该 topic 的订阅者。伪代码如下:
public interface ClusterPublisher { void publish(String topic, byte[] payload, MqttQoS qos, boolean retain); } public class RedisClusterPublisher implements ClusterPublisher { private final String channelName = "broker:route:event"; @Override public void publish(String topic, byte[] payload, MqttQoS qos, boolean retain) { byte[] event = eventCodec.encode(topic, payload, qos.value(), retain); redisTemplate.execute(conn -> conn.publish(channelName, event)); } }注意这里是异步还是同步调用。如果使用 Spring Data Redis 的模板同步执行 publish,会阻塞 Netty 的 EventLoop,量一大就出问题。生产环境建议用异步客户端;核心是不让集群广播的延迟影响到本地连接的读取。这种广播会带来消息重复:node1 发布,node1 自己也会收到一份,因此在接收端要判断消息来源节点,本机发布的消息不要重复投递。更严格的全网去重需要为每条消息分配全局唯一 ID,在投递端做幂等,否则订阅方可能收到两条相同消息。
4.3 用 RocksDB 持久化会话和 QoS1 待确认消息
集群节点可以重启,但业务侧不能把订阅关系全部丢掉。嵌入式数据库里我一般优先用 RocksDB,它不引入额外的服务端组件,部署时打包一个 native 库就能跑,适合 Java 系 broker。数据文件格式建议按用途分开,用 Column Family 区分 session、subscription、outgoing QoS1、retain message。写入时特别注意:不要在 EventLoop 线程里同步写 RocksDB,用一个单线程写队列聚合批量写,或者用 Netty 的业务线程池执行。代码层面对外暴露的保存逻辑可以是:
public void saveOutgoing(String clientId, MqttPublishMessage message) { byte[] key = ("qos1:" + clientId + ":" + message.variableHeader().packetId()).getBytes(StandardCharsets.UTF_8); byte[] value = encodeMessageWithHeaders(message); rocksDB.put(key, value); // 内部走异步写入队列,不在 EventLoop 上直接执行 }RocksDB 的参数不建议直接用默认值。需要结合内存和磁盘 IO 调整,常见参数如下:
| RocksDB 参数 | 建议值 | 用途 |
|---|---|---|
max_background_jobs | 4 到 8 | 控制 compaction 和 flush 线程,避免写放大抢 CPU |
write_buffer_size | 64MB 以上 | 减少小 value 写放大,按内存余量调整 |
level0_file_num_compaction_trigger | 4 | 避免 L0 文件过多导致读放大 |
max_open_files | 根据 fd 余量设置 | 控制 RocksDB 占用的文件描述符,10 万连接时 fd 很紧张 |
RocksDB 的数据本身可以定期备份到对象存储,也可以依赖多副本 backup。生产中单机进程挂掉后,重启打开同一份 RocksDB 数据就能把 QoS1 待确认消息重新加载,这比全内存方案可靠得多。
4.4 JVM 参数和 Netty 水位别抄默认值
10 万连接下,JVM 参数不能直接套 Web 应用的默认配置。连接本身在堆内占用不高,但 Netty 的堆外内存是按 channel 分配的,发送缓冲、接收缓冲都在堆外。如果不限制每条连接的 buffer,一个慢客户端就能把 Direct Memory 吃满。代码里设置WRITE_BUFFER_WATER_MARK只是第一步,JVM 层还要显式限制:
-Xms8g -Xmx8g -XX:MaxDirectMemorySize=1g -XX:+UseG1GC -XX:MaxGCPauseMillis=50 -XX:+ExitOnOutOfMemoryError-XX:MaxDirectMemorySize=1g不是越大越好,堆外内存超过物理内存后会在不可控的地方触发 OOM。G1 的停顿目标设到 50ms 是常见做法,但要注意 10 万连接时的 GC 日志里如果长期出现 humongous allocation,就要检查是否有人把大 byte[] 直接放进 EventLoop。老年代使用率涨上去不一定是泄漏,也可能是 QoS1 待确认消息全部积压在 ConcurrentHashMap 里,先看指标再调参数。
5. 生产验证:10 万连接的压测、系统参数和故障观测
到了验证阶段,代码层面的事基本收口,接下来全是操作系统和实验设计的问题。很多人把压测目标设成“能建立 10 万连接”,实际这是最基础的一步;连接起来后消息能不能稳定收发、broker 能否恢复,才是生产验证的重点。
5.1 先把连接数上限从系统层打开
在压测机上,先确认文件描述符上限。服务端每个连接至少一个 fd,加上 Netty 的 epoll、RocksDB 的文件句柄,建议ulimit -n至少 100 万。还有 TCP 端口范围:压测机作为客户端,发起的连接会占用本地端口,端口不够会报Cannot assign requested address。常见做法是:
ulimit -n 1048576 sysctl -w net.core.somaxconn=32768 sysctl -w net.ipv4.ip_local_port_range="1024 65535" sysctl -w net.ipv4.tcp_fin_timeout=30net.core.somaxconn要配合SO_BACKLOG,否则 accept 队列不够,大量并发建连时握手成功但连接建立缓慢。压测机上ip_local_port_range设置为 1024 到 65535,再结合连接复用,才能支持足够的源端口。tcp_fin_timeout影响 TIME_WAIT 回收速度,压测机端口不够时可以把默认 60 秒调低;broker 所在节点则不要随意开 tcp_tw_reuse,作为被动方它只影响主动出站连接,收益不大还可能有风险。
5.2 压测时该看的统计,而不是只看“连上了”
压测过程中,不要只看工具显示的“connected=100000”。直接在 broker 上运行ss -s,看系统当前 socket 总数和状态;用cat /proc/net/sockstat看 TCP 的 inuse、timewait 数量;然后用jstack看业务线程是否大量卡在同一个锁或同步 IO 上。如果ss -s显示连接数稳定,但jstack里 worker 线程全在RocksDB.put,说明持久化写进了 EventLoop,这是最典型的失败模式。
另一个容易忽略的指标是 channel 写水位。可以在自定义 handler 里定期采集 isWritable 状态,如果大量连接处于不可写,说明压测中的客户端消费速度跟不上 broker 下发速度,调大 worker 线程没用,应该检查订阅端的 TCP 窗口和 QoS1 积压队列。
5.3 消息计数校验:连接数不等于可用性
压测收尾前,我一般会跑三类验证:干净连接建立、心跳拨测、QoS1 消息计数校验。第三类最容易做错:只统计 broker 发出去多少,不统计客户端实际收到多少。正确做法是在发布端累计 PUBLISH 数量,在订阅端累计收到的 payload 序号,压测结束后比对两侧总数和遗漏区间。如果两端总数一致但中间有乱序,需要检查集群广播是否把消息从多个节点重复投递。连接数到 10 万只是入场券,消息计数对得上才敢把 broker 交出去。每次发布前我会在预发环境重复这套流程,任一项不过就继续查,直到三类任务都用同一个版本过掉。
本文还有配套的精品资源,点击获取