拓冰建站拓冰建站
首页 / 资讯中心 / 正文

基于Netty的高并发MQTT Broker实战:10万连接与生产级优化

简介压缩包内是一份基于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 只做 acceptworkerGroup 负责 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 ChannelInitializerSocketChannel() { 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换成EpollEventLoopGroupchannel 对应换成EpollServerSocketChannel能减少一次从 select 到 epoll 的系统调用开销。SO_BACKLOG是 accept 队列长度只写 1024 还不够后面要配合操作系统参数一起看。常见参数可按下表设置参数建议值说明SO_BACKLOG1024 或 2048TCP accept 队列长度超过net.core.somaxconn时以内核为准SO_REUSEADDRtrue快速重启 broker避免端口处于 TIME_WAIT 时无法 bindTCP_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就是第一道闸门当一条连接未发送的字节数超过高水位 1MBchannel 变成不可写业务代码需要停止继续投递等它降到低水位以下再继续。实现上需要检查channel.isWritable()不要无脑writeAndFlushif (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 ConcurrentHashMapString, ChannelHandlerContext channelMap new ConcurrentHashMap(); private final ConcurrentHashMapString, 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 后删除重发时带 DUPQoS 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 线程避免写放大抢 CPUwrite_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:MaxDirectMemorySize1g -XX:UseG1GC -XX:MaxGCPauseMillis50 -XX:ExitOnOutOfMemoryError-XX:MaxDirectMemorySize1g不是越大越好堆外内存超过物理内存后会在不可控的地方触发 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.somaxconn32768 sysctl -w net.ipv4.ip_local_port_range1024 65535 sysctl -w net.ipv4.tcp_fin_timeout30net.core.somaxconn要配合SO_BACKLOG否则 accept 队列不够大量并发建连时握手成功但连接建立缓慢。压测机上ip_local_port_range设置为 1024 到 65535再结合连接复用才能支持足够的源端口。tcp_fin_timeout影响 TIME_WAIT 回收速度压测机端口不够时可以把默认 60 秒调低broker 所在节点则不要随意开 tcp_tw_reuse作为被动方它只影响主动出站连接收益不大还可能有风险。5.2 压测时该看的统计而不是只看“连上了”压测过程中不要只看工具显示的“connected100000”。直接在 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 交出去。每次发布前我会在预发环境重复这套流程任一项不过就继续查直到三类任务都用同一个版本过掉。本文还有配套的精品资源点击获取
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门