1. 为什么我放弃了Spring原生WebSocket,转头上了Netty
先交代一下背景。我这边有个项目,前期用的是SpringBoot自带的WebSocket(基于WebSocketHandler和STOMP那套),单机几百个连接的时候一切正常。等业务量上来,在线连接冲到两三千,问题就开始冒头了:时不时有连接假死、心跳超时误判、偶尔出现消息延迟,更离谱的是Tomcat容器偶尔直接报线程池耗尽。
排查了一圈,根子不在业务代码,而在协议栈和线程模型上。Spring的WebSocket底层还是走Servlet容器那一套,一个连接一个请求处理的思路,长连接场景下线程占用非常不划算。后来我把WebSocket这块整个切到了Netty,连接数到一万级别依然稳定。今天这篇就把整个接入过程、代码骨架、以及我当时踩过的坑完整写一遍,包含粘包拆包、心跳保活、Nginx反代这些生产环境绕不开的细节。
“三分钟构建”不是说三分钟把所有业务写完,而是指:接入链路足够清晰,从零到跑通一条可用连接,整个核心代码量控制在一个合理范围内,剩下的时间都花在业务处理和调优上。如果你正打算在SpringBoot项目里上WebSocket,又担心原生方案撑不住以后的高并发,这篇应该能帮你少走不少弯路。
适用人群分两类:一类是刚接触Netty、想在SpringBoot里跑通WebSocket的开发者;另一类是已经在用原生WebSocket、但被连接数、稳定性问题折腾过的人。前者重点看第二章和第三章的代码骨架,后者重点看第四章和第五章的踩坑记录与保活策略。
2. 方案选型不纠结:Netty到底解决了原生方案的哪些问题
很多人在选型时会犹豫:SpringBoot官方就支持WebSocket,为什么要额外引入Netty?这里我不打算只给结论,把两者在真实运行时的差异拆开讲。
2.1 原生WebSocket的瓶颈不在协议,在线程模型
Spring的原生WebSocket实现(包括STOMP)本质上是建立在Servlet容器之上的。Servlet 3.1之后虽然支持异步处理,但长连接场景下每个WebSocket连接依然对应着一个独立的“逻辑处理单元”,高并发时容器线程池很容易成为瓶颈。
用大白话说:Tomcat处理普通HTTP请求是“来一个请求、分配一个线程、处理完归还”,这种模式对短连接没问题。但WebSocket是长连接,连接建立后通道一直占用着,如果线程模型设计不合理,大量空闲连接也会把线程池资源耗尽,导致正常的HTTP请求都进不来。
我遇到的具体现象就是:在线连接数到一定量级后,登录接口、业务查询接口这些普通HTTP请求开始变慢,甚至偶发超时。压测时看监控,容器线程池活跃数长期处于高位。这就是典型的“长连接吃掉了短连接的资源”。
2.2 Netty的NIO模型为什么适合长连接场景
Netty基于NIO(非阻塞I/O)事件驱动模型,核心是EventLoop线程组。它和Servlet容器的关键差异在于:一个EventLoop线程可以同时管理成千上万个Channel(连接),而不是一个连接一个线程。连接上的读写事件来了才触发处理,没有事件时线程可以去处理其他连接的事件。
这就好比餐厅服务员:原生方案是一桌客人配一个专属服务员,哪怕客人只是坐着聊天,服务员也得在旁边候着;Netty是一组服务员轮流巡视所有桌子,只有客人真正举手示意(有事件)时才过去服务。同样的人力,后者能服务的桌数要远多于前者。
这也是为什么高并发长连接场景下,Netty的线程占用更少、支撑的连接数更多。WebSocket协议本身在传输层就是基于TCP的,而Netty在TCP层面做了大量优化,比如零拷贝、内存池、更细粒度的读写控制,这些对追求极致性能的场景都非常有价值。
2.3 方案边界:什么情况下其实不必上Netty
这里说点实话,不是所有项目都需要上Netty。如果你的业务是内部管理系统、运维后台,在线连接数长期只有几十到几百,用Spring原生WebSocket完全够用,引入Netty反而增加维护成本——毕竟多了一套线程模型、多了一种协议处理方式,团队学习成本是实打实的。
我的判断标准很简单:
| 判断维度 | 适用原生WebSocket | 适用SpringBoot + Netty |
|---|---|---|
| 在线连接数 | 百级以内 | 千级以上或预期快速增长 |
| 消息频率 | 低频(秒级心跳、少量推送) | 高频(毫秒级推送、实时交互) |
| 线程模型要求 | 无特殊要求 | 需要控制线程数、支撑海量连接 |
| 协议定制需求 | 无,标准文本/二进制帧即可 | 需要拆包、粘包处理、自定义帧协商 |
| 团队技术储备 | Spring技术栈熟练 | 愿意投入少量时间学习Netty核心概念 |
我自己现在的做法是:核心消息网关用Netty;周边简单工具页面、内部通知类的小服务用Spring原生WebSocket,怎么省事怎么来。技术选型没有绝对的对错,只有合适不合适。
3. 三分钟快速搭建:完整的SpringBoot + Netty接口骨架
下面进入正题,直接给可复制的代码。我用的环境是SpringBoot 2.7.x + Netty 4.1.x,这两个版本是目前生产环境里最稳的组合,JDK 8及以上都能跑。
3.1 Maven依赖引入
<dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> <dependency> <groupId>io.netty</groupId> <artifactId>netty-all</artifactId> <version>4.1.100.Final</version> </dependency>很多教程会建议直接在SpringBoot项目里引入netty-all,实际开发中也可以按需引入netty-transport、netty-codec-http、netty-handler等模块。不过netty-all省心,不差那点依赖体积的话直接用也未尝不可。
3.2 Netty服务启动类
先定义一个组件,负责在SpringBoot启动完成后拉起Netty服务端。这里要注意的是:不能直接在SpringBoot的main方法里启动Netty,因为那会导致Netty和Spring容器的生命周期不一致。我用的是ApplicationRunner,Spring容器初始化完成后自动执行。
@Component public class NettyWebSocketServer implements ApplicationRunner { private static final Logger log = LoggerFactory.getLogger(NettyWebSocketServer.class); private EventLoopGroup bossGroup; private EventLoopGroup workerGroup; private Channel serverChannel; private final WebSocketServerInitializer initializer; public NettyWebSocketServer(WebSocketServerInitializer initializer) { this.initializer = initializer; } @Override public void run(ApplicationArguments args) { bossGroup = new NioEventLoopGroup(1); workerGroup = new NioEventLoopGroup(); try { ServerBootstrap bootstrap = new ServerBootstrap(); bootstrap.group(bossGroup, workerGroup) .channel(NioServerSocketChannel.class) .option(ChannelOption.SO_BACKLOG, 1024) .childOption(ChannelOption.TCP_NODELAY, true) .childOption(ChannelOption.SO_KEEPALIVE, true) .childHandler(initializer); int port = 8899; ChannelFuture future = bootstrap.bind(port).sync(); serverChannel = future.channel(); log.info("Netty WebSocket服务启动成功,端口:{}", port); } catch (InterruptedException e) { Thread.currentThread().interrupt(); log.error("Netty WebSocket服务启动失败", e); } } @PreDestroy public void destroy() { if (serverChannel != null) { serverChannel.close(); } if (bossGroup != null) { bossGroup.shutdownGracefully(); } if (workerGroup != null) { workerGroup.shutdownGracefully(); } } }这里有几个参数值得解释一下。bossGroup线程数设为1就够了,因为它只负责接受新的连接请求,把连接注册到workerGroup。workerGroup不传线程数时,默认是CPU核数的两倍,这个配置在绝大多数场景下都是合理的。TCP_NODELAY设为true是为了禁用Nagle算法,避免小数据包被延迟合并发送——WebSocket是实时交互协议,延迟比带宽更敏感。
顺带说一句,关于SpringBoot版本的问题,很多人在热搜里提到“SpringBoot版本太高不能使用JDK1.8”。这是因为SpringBoot 3.x开始强制要求JDK 17,如果你的生产环境还在用JDK 8,那最好固定使用SpringBoot 2.7.x版本;反过来,如果你已经用JDK 17,那可以用SpringBoot 3.x,Netty本身对JDK版本没有强限制,兼容性反而比SpringBoot更宽。我这边生产环境是从SpringBoot 2.x起步的,后续要升级到3.x也得先把JDK升上去,这条链路建议在项目初期就想清楚。
3.3 ChannelInitializer:Pipeline里各Handler的装配顺序
Netty里最核心的概念之一,就是每个连接的Pipeline,也就是责任链。Handler的添加顺序很关键,因为数据帧在链路上是顺序流过的。
@Component public class WebSocketServerInitializer extends ChannelInitializer<SocketChannel> { private final WebSocketFrameHandler webSocketFrameHandler; public WebSocketServerInitializer(WebSocketFrameHandler webSocketFrameHandler) { this.webSocketFrameHandler = webSocketFrameHandler; } @Override protected void initChannel(SocketChannel ch) { ChannelPipeline pipeline = ch.pipeline(); // 1. 粘包/半包处理 pipeline.addLast(new LengthFieldBasedFrameDecoder(1024 * 1024, 0, 4, 0, 4)); // 2. HTTP协议编解码 pipeline.addLast(new HttpServerCodec()); // 3. 聚合HTTP请求 pipeline.addLast(new HttpObjectAggregator(65536)); // 4. WebSocket协议升级 pipeline.addLast(new WebSocketServerProtocolHandler("/ws")); // 5. 业务处理 pipeline.addLast(webSocketFrameHandler); } }注意,这里我在HttpServerCodec之前加了一个LengthFieldBasedFrameDecoder,这是为了处理粘包/半包问题。这个话题比较重要,第三章会专门展开。需要先说明的是,WebSocket协议本身已经有帧边界了,但底层的TCP流是字节流,如果前置的编解码器不能正确识别帧边界,就会把多个帧混在一起交给后面的处理器,导致解析出错。这个坑很隐蔽,很多人写了半天总发现消息对不上,最后定位到就是这一步的问题。
WebSocketServerProtocolHandler("/ws")表示WebSocket的路径是/ws,客户端连接时要连ws://ip:8899/ws。如果你要支持多个路径,就需要多注册几个Handler,或者自己在业务Handler里做路由判断。
3.4 业务Handler:接收消息、推送消息、连接管理
这个是核心处理类,实现了连接建立、消息收发、连接关闭的主要逻辑。
@Component @ChannelHandler.Sharable public class WebSocketFrameHandler extends SimpleChannelInboundHandler<WebSocketFrame> { private static final Logger log = LoggerFactory.getLogger(WebSocketFrameHandler.class); /** * 管理所有在线连接,concurrent包做线程安全 */ private static final ChannelGroup onlineChannels = new DefaultChannelGroup(GlobalEventExecutor.INSTANCE); /** * 记录ChannelId和用户ID的映射,方便定向推送 */ private static final ConcurrentHashMap<String, String> userChannelMap = new ConcurrentHashMap<>(); @Override public void channelActive(ChannelHandlerContext ctx) { onlineChannels.add(ctx.channel()); log.info("连接建立:{},当前在线连接数:{}", ctx.channel().remoteAddress(), onlineChannels.size()); } @Override public void channelInactive(ChannelHandlerContext ctx) { onlineChannels.remove(ctx.channel()); userChannelMap.entrySet().removeIf(entry -> entry.getValue().equals(ctx.channel().id().asLongText())); log.info("连接断开:{},当前在线连接数:{}", ctx.channel().remoteAddress(), onlineChannels.size()); } @Override protected void channelRead0(ChannelHandlerContext ctx, WebSocketFrame frame) { // 关闭帧 if (frame instanceof CloseWebSocketFrame) { ctx.close(); return; } // Ping/Pong帧 if (frame instanceof PingWebSocketFrame) { ctx.channel().writeAndFlush(new PongWebSocketFrame(frame.content().retain())); return; } if (frame instanceof TextWebSocketFrame) { String message = ((TextWebSocketFrame) frame).text(); log.info("收到消息:{}", message); // 这里就可以根据消息内容做业务分发了 // 例如:用户认证、指令下发、广播等 handleTextMessage(ctx, message); return; } if (frame instanceof BinaryWebSocketFrame) { // 处理二进制数据 log.info("收到二进制数据,长度:{}", frame.content().readableBytes()); return; } } private void handleTextMessage(ChannelHandlerContext ctx, String message) { // 一个简单的消息协议,例如:{"type":"auth","userId":"10001"} try { JSONObject json = JSON.parseObject(message); String type = json.getString("type"); if ("auth".equals(type)) { String userId = json.getString("userId"); userChannelMap.put(userId, ctx.channel().id().asLongText()); ctx.channel().writeAndFlush(new TextWebSocketFrame("认证成功")); } } catch (Exception e) { ctx.channel().writeAndFlush(new TextWebSocketFrame("消息格式错误")); } } @Override public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { log.error("WebSocket连接异常", cause); ctx.close(); } /** * 向指定用户推送消息 */ public static boolean sendToUser(String userId, String message) { String channelId = userChannelMap.get(userId); if (channelId == null) { return false; } for (Channel channel : onlineChannels) { if (channel.id().asLongText().equals(channelId)) { channel.writeAndFlush(new TextWebSocketFrame(message)); return true; } } return false; } /** * 广播消息给所有在线连接 */ public static void broadcast(String message) { onlineChannels.writeAndFlush(new TextWebSocketFrame(message)); } }这里使用了ChannelGroup来管理在线连接,好处是它的writeAndFlush可以批量向所有Channel发消息,而且线程安全。SimpleChannelInboundHandler会自动释放ReferenceCounted对象,所以处理完消息后不需要手动release,省去不少内存泄漏隐患。
还有一个细节:@ChannelHandler.Sharable这个注解一定要加,因为我把Handler作为Spring单例Bean放进了Pipeline,多个Channel共享同一个Handler实例。如果你不确信自己写的是无状态的,那就别加这个注解,每次初始化时new一个Handler更保险。我的这个Handler里所有状态都在static字段上,所以是安全的。
3.5 运行效果验证
把项目跑起来,控制台会输出Netty启动成功日志。再用Postman连一下测试链路是否通。
Postman从某个版本开始就原生支持WebSocket请求了,打开Postman,左侧选择New,然后选WebSocket Request,输入ws://localhost:8899/ws,点击Connect。连接成功后,在Message输入框输入任意文本,点击Send,应该能收到服务端返回的消息。断开连接时,控制台会打印连接断开日志,在线连接数同步减少。
这个操作流程也对应了热搜里“postman websocket连接”这个需求,实测下来Postman对WebSocket的调试支持足够日常开发用了,不必专门找WebSocket调试工具。
4. 粘包与拆包:长连接里最容易翻车的环节
既然搜热词里专门有“netty粘包处理”,这块我必须单独拿出来详细讲。它也是新手从HTTP转到Netty后最先遇到的坑,没有之一。
4.1 粘包到底是怎么产生的
先说结论:TCP本身就是字节流协议,它不关心你上层发的是几条消息。TCP为了保证传输效率,会根据MSS(最大报文段大小)和当前网络缓冲区情况,把多个小数据包合并成一个TCP报文发送,或者把一个大数据包拆成多个TCP报文发送。这就是“粘包”和“半包”的来源。
粘包:发送方发送了两条消息hello和world,接收方一次读到了helloworld,无法确定消息边界。 半包:发送方发送了一条很长的消息,接收方一次只读到了前半段,需要等后续数据到达才能凑齐一条完整消息。
WebSocket协议号称自带消息边界,为什么还会遇到粘包?因为WebSocket的帧边界信息是在帧头(Frame Header)里声明的,如果解码器不能正确读取帧头信息,或者干脆把WebSocket帧当成裸TCP数据流来处理,就会出错。更常见的是在加WebSocketServerProtocolHandler之前,如果前面没有正确的拆包器,多个WebSocket帧的数据会先后到达,后面的处理器无法区分帧边界。
4.2 四种常见拆包器的选择
Netty提供了多种拆包器,适合不同的场景:
| 拆包器 | 原理 | 适用场景 | 注意点 |
|---|---|---|---|
LineBasedFrameDecoder | 按换行符\n或\r\n拆包 | 文本协议,每条消息以换行结束 | 消息内容不能包含换行符 |
DelimiterBasedFrameDecoder | 按自定义分隔符拆包 | 自定义文本协议 | 分隔符不能出现在业务数据内 |
FixedLengthFrameDecoder | 按固定长度拆包 | 消息定长 | 不易扩展,不推荐业务用 |
LengthFieldBasedFrameDecoder | 按长度字段拆包 | 二进制协议、混合协议 | 参数配置较复杂,但最通用 |
WebSocket场景下,最稳的方案是在HTTP协议处理器前加LengthFieldBasedFrameDecoder,并且正确配置长度字段的位置。但要注意:WebSocket的帧头和普通的TCP长度字段协议不一样,如果只是加一个LengthFieldBasedFrameDecoder而不告诉它WebSocket帧的头部结构,那拆出来的仍然会是乱的。
更准确的做法是:如果你确定走的是WebSocket协议升级,可以在升级完成后依靠WebSocketServerProtocolHandler来维护WebSocket帧的边界,这个Handler内部实现了对WebSocket帧的自动拆分,不需要额外处理。但问题在于,在握手之前,连接还是普通的HTTP连接,此时如果数据帧到达,必须先经过HTTP解码器,而这个阶段的粘包处理需要HttpObjectAggregator配合。
我个人的实践中,最省心的排列方式是:
pipeline.addLast(new HttpServerCodec()); pipeline.addLast(new HttpObjectAggregator(65536)); pipeline.addLast(new WebSocketServerProtocolHandler("/ws"));在这个顺序下,HttpObjectAggregator会聚合HTTP请求和升级握手请求,WebSocketServerProtocolHandler在升级完成后接管帧解析。如果业务消息需要走自定义二进制协议,或者想绕开WebSocket协议、直接做TCP长连接,这个时候LengthFieldBasedFrameDecoder才会成为主力。
4.3 一个真实现场的粘包复现
我之前压测时遇到过这个情况:客户端大批量推送消息,服务端偶尔出现解析出来的消息是两条拼在一起的现象。当时监控看错误日志,发现收到的消息末尾会莫名多出一段上一个消息的内容。
这就是典型的粘包。定位时打了不少日志,发现发生粘包的位置集中在连接刚建立的瞬间,或者消息体较大的场景。原因很简单:发送方在两毫秒内连发了多条消息,TCP层的Nagle算法和接收方的缓冲区机制把这些消息合并成了同一个TCP报文。如果解码器没有按帧边界正确切分,业务层拿到的就是“连体婴儿”。
解决思路有两种:
第一种,如果是WebSocket协议,确保WebSocketServerProtocolHandler之后不要再加容易破坏帧边界的Handler,并确认业务Handler直接接收的是WebSocketFrame子类。
第二种,如果是自己定义的二进制协议,用LengthFieldBasedFrameDecoder明确告诉Netty:前4个字节是消息总长度,收到这个长度后再交给后续Handler。配置代码如下:
pipeline.addLast(new LengthFieldBasedFrameDecoder(1024 * 1024, 0, 4, 0, 4));这里的参数含义是:最大帧长度为1MB,长度字段偏移量为0,长度字段长度为4字节,长度字段后续还有0字节的头部数据,数据长度加上头部的调整量为4字节(即前4个字节代表后续数据的长度)。这四个参数如果不对,会直接导致解析错乱,建议对照你自己的协议头结构写清楚注释。
4.4 排查粘包问题的完整思路
万一你线上还是出现了粘包问题,排查路径按照这个顺序走,能少走弯路:
- 先确认协议版本和Handler顺序,把Pipeline里每个Handler的职责用注释标出来。
- 查看日志,确定触发粘包时客户端连发了几条消息、间隔是多大。
- 在拆包器前后各加日志,打印
ByteBuf的可读字节数和读取的内容,确认拆包器是否正常工作。 - 用小包连发、大包单发、大包连发三组用例分别测试,看问题出现在哪一类。
- 确认
Content-Length头是否正确(HTTP阶段)或WebSocket帧长度字段是否正确(升级后阶段)。
我之前就是靠第3步找到问题的:日志显示拆包器输出的一帧里包含了两个业务消息,说明拆包器根本没工作,后来发现是Pipeline里加的顺序不对,LengthFieldBasedFrameDecoder加在了HttpServerCodec之后,它拿到的数据已经被HTTP解码器吃掉了。
5. 心跳保活与断线重连:长连接稳定性的基石
WebSocket服务上线后,第二个高频问题就是连接“假死”——客户端以为还连着,服务端也还不知道连接已经断了,但消息发过去没有任何响应。这通常是网络中任何一环(尤其是Nginx或云负载均衡)把空闲连接静默关掉了。要解决这个问题,必须做心跳保活。
5.1 服务端空闲检测机制
Netty里最常用的心跳方案是IdleStateHandler,它可以在连接空闲超过指定时间后触发事件。我在Pipeline里加了一个专门负责心跳检测的Handler,放在业务Handler的最前面。
public class HeartbeatHandler extends ChannelInboundHandlerAdapter { private static final Logger log = LoggerFactory.getLogger(HeartbeatHandler.class); /** 配置参数:读空闲时间、写空闲时间、总空闲时间 */ private static final int READER_IDLE_TIME = 60; private static final int WRITER_IDLE_TIME = 30; private static final int ALL_IDLE_TIME = 0; @Override public void handlerAdded(ChannelHandlerContext ctx) { ctx.pipeline().addFirst(new IdleStateHandler(READER_IDLE_TIME, WRITER_IDLE_TIME, ALL_IDLE_TIME, TimeUnit.SECONDS)); } @Override public void userEventTriggered(ChannelHandlerContext ctx, Object evt) throws Exception { if (evt instanceof IdleStateEvent) { IdleStateEvent event = (IdleStateEvent) evt; if (event.state() == IdleState.READER_IDLE) { // 读空闲:客户端超过60秒没发任何数据,主动断开 log.info("{} 读空闲超时,关闭连接", ctx.channel().remoteAddress()); ctx.close(); } else if (event.state() == IdleState.WRITER_IDLE) { // 写空闲:向客户端发送心跳包 ctx.channel().writeAndFlush(new TextWebSocketFrame("{\"type\":\"heartbeat\"}")); } } super.userEventTriggered(ctx, evt); } }这里的设计思路是:服务端每30秒主动发一次心跳(写空闲触发),如果连续60秒没有读到客户端任何数据(读空闲触发),就认为客户端已经不活跃,主动关闭连接。这样的好处是能及时清理死连接,避免连接数被无效连接慢慢占满。
5.2 客户端的断线重连策略
服务端只是做被动清理,真正保证连接稳定的是客户端要有断线重连能力。很多前端方案里,WebSocket对象一断就什么都不管了,用户必须刷新页面才能恢复,体验很差。
前端的断线重连建议采用“指数退避”策略:第一次失败后等1秒再重连,第二次等2秒,第三次等4秒,以此类推,到最大间隔后封顶。并且要加上随机抖动,防止大量客户端同时重连造成服务端瞬间压力过大。
一个简化版的前端重连逻辑:
let reconnectAttempts = 0; const MAX_RECONNECT_DELAY = 30000; // 最大30秒 const BASE_DELAY = 1000; function connectWebSocket() { const ws = new WebSocket('ws://localhost:8899/ws'); ws.onopen = () => { reconnectAttempts = 0; console.log('WebSocket连接已建立'); // 连接成功后开始心跳 startHeartbeat(); }; ws.onclose = () => { clearHeartbeat(); const delay = Math.min(BASE_DELAY * Math.pow(2, reconnectAttempts), MAX_RECONNECT_DELAY) + Math.random() * 1000; reconnectAttempts++; console.log(`连接断开,${delay}ms后重连`); setTimeout(connectWebSocket, delay); }; ws.onerror = (err) => { console.error('WebSocket错误', err); ws.close(); }; } // 客户端心跳,每25秒发送一次,比服务端的30秒略短,确保连接活跃 function startHeartbeat() { window.heartbeatTimer = setInterval(() => { if (ws.readyState === WebSocket.OPEN) { ws.send('{"type":"heartbeat"}'); } }, 25000); }这个策略在真实项目里很管用。客户端心跳间隔比服务端空闲阈值略短,能保证连接在服务端判定超时之前就有数据流动,避免被误杀。服务端收到心跳消息后,可以回一个pong消息,也可以不回——因为收到消息这个动作本身就刷新了读空闲检测的计时。
5.3 服务端踢人策略
除了心跳,还有一类场景:同一账号在多个终端登录,或者恶意客户端建立了一堆连接不干正事。这时候服务端需要有主动踢人的能力。
我的做法是,在认证通过后建立userId -> Channel的映射,如果同一个userId已经有映射,就主动关闭旧连接,让新连接接管。这里要注意并发问题:两个连接同时认证同一个userId,可能互相踢,所以需要用ConcurrentHashMap的putIfAbsent或者加锁来保证原子性。
简化逻辑如下:
String oldChannelId = userChannelMap.put(userId, ctx.channel().id().asLongText()); if (oldChannelId != null && !oldChannelId.equals(ctx.channel().id().asLongText())) { for (Channel ch : onlineChannels) { if (ch.id().asLongText().equals(oldChannelId)) { ch.writeAndFlush(new TextWebSocketFrame("{\"type\":\"kick\",\"reason\":\"account_login_elsewhere\"}")); ch.close(); break; } } }这个机制还有个好处:配合断线重连,客户端重连时旧连接会被自动清理,不会出现“连接残留”导致消息发给一个已经废弃的Channel。
6. Nginx反代、压测数据与生产环境避坑指南
WebSocket在生产环境通常不会让客户端直连应用服务,前面会挡一层Nginx做负载均衡和域名分发。但Nginx默认配置不支持WebSocket,必须显式开启升级头。
6.1 Nginx配置WebSocket的关键项
map $http_upgrade $connection_upgrade { default upgrade; '' close; } upstream websocket_backend { server 127.0.0.1:8899; # 如果有多台服务,可以加权重 # server 127.0.0.1:8900 weight=2; keepalive 32; } server { listen 80; server_name ws.example.com; location /ws { proxy_pass http://websocket_backend; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection $connection_upgrade; proxy_set_header Host $host; proxy_set_header X-Real-IP $remote_addr; proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for; proxy_read_timeout 3600s; proxy_send_timeout 3600s; } }map块的两行配置很关键,它告诉Nginx:如果客户端请求头里有Upgrade: websocket,就把Connection头设为upgrade,否则保持默认的close。proxy_read_timeout默认是60秒,如果不改,WebSocket长连接空闲超过60秒就会被Nginx掐断——很多开发环境连得好好的,一上Nginx就频繁掉线,十有八九就是这个超时时间没调大。
我这边把超时设为3600秒(1小时),实际运营中客户端和服务端有心跳保活,Nginx层面的空闲其实很少发生,但设长一点可以避免极端情况下的误断。
6.2 压测结果对比
我用同一台4核8G的云服务器做过一轮简单压测,分别压Spring原生WebSocket和SpringBoot + Netty方案。压测工具是JMeter加WebSocket插件,模拟客户端持续连接和发送消息。结果数据因为测试环境有波动,但趋势非常一致:
| 指标 | Spring原生WebSocket | SpringBoot + Netty |
|---|---|---|
| 5000并发连接 | 内存占用持续上涨,部分连接超时 | 稳定运行,内存平稳 |
| 10000并发连接 | 连接大量失败,服务日志报线程池满 | 稳定运行,CPU约40% |
| 消息推送延迟(P99) | 高负载时偶发秒级延迟 | 稳定在几十毫秒内 |
| 空闲连接占用资源 | 每连接一个线程,资源占用大 | 每连接一个Channel,事件驱动,资源占用小 |
这个压测严格说不够严谨,毕竟压测期间网络、机器负载都会有波动,但足以说明方向:在连接数和消息量上来之后,Netty的线程模型优势非常明显。如果你的项目预见到未来在线连接数会超过几千,尽早切到Netty是划算的。
6.3 容易被忽略的内存泄漏和异常处理
生产环境跑了一段时间后,我遇到过两类问题,都值得单独提醒。
第一类是ByteBuf泄漏。Netty使用堆外内存的ByteBuf时,必须手动release,忘了释放就会导致堆外内存泄漏,最终表现为内存占用只涨不降、甚至OOM。解决办法是:用好SimpleChannelInboundHandler(它自动释放),如果继承的是普通ChannelInboundHandlerAdapter,要在channelRead处理完手动ReferenceCountUtil.release(msg),或者ctx.writeAndFlush会自动释放。另外,设置-Dio.netty.leakDetection.level=advanced可以在日志里打印泄漏点,定位时建议开起来。
第二类是连接数监控缺失。线上如果连接数异常暴涨,却没有监控报警,等发现问题时服务可能已经不可用了。我在服务端加了一个定时任务,每10分钟输出一次当前在线连接数和ChannelGroup大小,配合Prometheus暴露指标,做到趋势可见。
@Scheduled(fixedRate = 600000) public void reportConnectionCount() { log.info("当前在线连接数:{},UserChannel映射数:{}", onlineChannels.size(), userChannelMap.size()); }这个日志看似简单,关键时刻能救命。例如某次线上出现大量僵尸连接,就是靠这个日志里的映射数远大于正常业务用户数才发现的——最终定位到是某类客户端断网后没有自动重连也没有正常关闭连接,服务端靠心跳也清不掉(因为还在收心跳),最后靠限制单用户最大连接数和更严格的空闲清理解决。
再补充一点关于userChannelMap映射清理的心得:很多教程只做channelInactive时移除映射,但实际生产里,客户端进程被直接杀掉、网络闪断等场景,服务端不一定能及时触发channelInactive。所以要在心跳Handler里加上对映射的检查——如果连接被判定为读空闲并关闭,一定要把对应的映射也删掉,否则userChannelMap会残留大量废弃映射,定向推送时白白遍历一堆无效Channel。
7. 我给团队定下的编码规范,附带写给你的一些建议
项目稳定运行几个月后,我梳理了一套适合团队新人快速上手的编码规范,顺手分享几个最有价值的小技巧。
第一,所有发出去的消息都统一走封装方法,不要直接writeAndFlush。我定义了一个MessageSender工具类,里面封装了sendToUser、broadcast、sendToChannel等方法,统一处理消息式的JSON序列化、异常捕获、发送回调。好处是后续如果要加消息轨迹、限流、敏感词过滤,只需要改一个入口。
第二,WebSocketSession和业务线程之间的消息推送要串行化处理。如果业务线程向同一个Channel高频推送消息,不控制并发的话,writeAndFlush可能会出现线程安全问题。Netty本身保证了一个Channel上的操作是串行的,但前提是你不能从多个线程同时写同一个Channel而没有同步。我这边用一个ChannelFutureListener在写失败时打日志,同时避免在业务线程里阻塞等待写结果。
第三,Client与Server之间的消息协议一定要带type字段和requestId字段。刚开始我只用了type字段区分业务类型,后来做日志追踪和去重时发现,没有requestId几乎没法做联调。加一个UUID或者业务流水号,几行改动,排查问题的时间能省一大截。
第四,从SpringBoot 2.x升级到3.x时,Spring的WebMvcConfigurer接口和很多Web相关的自动配置都变了,如果项目里还依赖了其他Web框架,升级要通盘考虑,不只是换JDK版本那么简单。Netty和SpringBoot的集成代码在3.x下倒是基本不变,因为Netty自己管理生命周期,与Spring容器只有极少的耦合点。
我在实际项目中还有一个习惯,会把Netty端口注册到Nacos配置中心,而不是硬编码在代码里。这样做的好处是,环境切换(开发、测试、生产)时不必改代码重新打包,只要改配置就行。如果你项目里已经用Nacos或者配置中心,强烈建议这么干。
最后说点掏心窝的话:Netty这套方案好归好,但不要一上来就追求极致的性能调优。先把连接管理、心跳保活、粘包处理这几个基础环节做扎实,性能自然不会差到哪里去。我见过不少团队一上来就调ByteBuf池参数、折腾EventLoop线程数,结果基础业务还没跑通,出了问题也排除不了。先把正确的事情做对,再去优化跑得快不快。