基于Netty的高并发MQTT Broker实战:10万连接与生产级优化
2026/9/13 21:10:42 网站建设 项目流程

简介:压缩包内是一份基于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_BACKLOG1024 或 2048TCP accept 队列长度,超过net.core.somaxconn时以内核为准
SO_REUSEADDRtrue快速重启 broker,避免端口处于 TIME_WAIT 时无法 bind
TCP_NODELAYtrueMQTT 报文通常很小,禁用 Nagle 可减少 40ms 延迟
SO_KEEPALIVEtrue先交给内核探测死链,业务层心跳继续自己处理
WRITE_BUFFER_WATER_MARK512KB / 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_jobs4 到 8控制 compaction 和 flush 线程,避免写放大抢 CPU
write_buffer_size64MB 以上减少小 value 写放大,按内存余量调整
level0_file_num_compaction_trigger4避免 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=30

net.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 交出去。每次发布前我会在预发环境重复这套流程,任一项不过就继续查,直到三类任务都用同一个版本过掉。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询