尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

SpringBoot整合Netty WebSocket:高并发消息推送实战

SpringBoot整合Netty WebSocket:高并发消息推送实战 简介面向 Java 后端开发者的一份 SpringBootWebSocketNetty 消息推送实战示例文档适合已掌握 Spring 基础、希望补齐实时通信与长连接开发能力的读者参考可用于聊天、行情刷新、系统通知等需要服务端主动推送数据的场景。内容从依赖配置讲起逐步展开 NettyConfig 中 ChannelGroup 对全部连接的管理、ConcurrentHashMap 维护用户 ID 与 Channel 映射的做法NettyServer 内 bossGroup 与 workerGroup 的分工及独立线程启动方式以及 WebSocketHandler 继承 SimpleChannelInboundHandler 后对连接建立、消息接收与关闭事件的处理前端则用 JavaScript 的 WebSocket API 建连并监听 message 事件把推送内容渲染到页面文本域。读者可据此实现向全体客户端广播、按用户 ID 定向推送两类逻辑代码注释完整便于对照调试与二次改造。包内共 1 个 PDF 文件约 172KB篇幅紧凑可一次读完。已有 6276 人学习适合需要快速搭建实时消息推送原型的中高级开发者参考。1. 消息推送选型先想清楚什么时候 SpringBoot 内置 WebSocket 不够需要把 Netty 拉进来很多人第一次做消息推送是从 SpringBoot 的 SseEmitter 或者 spring-boot-starter-websocket 起步的几十上百个在线连接时看不出问题等连接数爬到几千线程栈、GC 停顿和慢客户端就会把整个服务拖住。标题里把 SpringBoot、WebSocket、Netty 三个词放在一起讲的其实是一件事的两种分工Netty 专职管连接、编解码、心跳和字节搬运SpringBoot 专职管依赖注入、事务、Redis、数据库和会话路由两者跑在同一个 JVM 进程里通过一张内存路由表对接。这套结构适合后台系统的站内消息推送、工单与审批提醒、IM 雏形、设备指令下发这类需要长连接在线推送的场景如果一天只推几条通知、在线连接长期在三位数以内用 Spring 自带的 WebSocket 或者干脆用 SSE 更省事。下面按选型理由、工程搭建、推送链路、粘包与调优的顺序,把能直接抄的示例代码和参数逐个讲清楚。2. WebSocket 服务端的分层设计SpringBoot 管业务Netty 管连接2.1 三条候选路线SSE、Spring 内置 WebSocket、Netty WebSocket先说清楚 sse消息推送是什么意思SSE 是服务器单向推送走 HTTP 长连接浏览器侧用 EventSource 接收断线会自动重连实现成本最低但它只支持服务端到客户端一个方向客户端想上行还得另开接口。Spring 内置的 WebSocketServerEndpoint或WebSocketHandler跑在 Servlet 容器里握手、缓冲区、线程调度都要看容器的脸色Tomcat 换 Undertow 行为就不一样。Netty 则是自己完全掌控 pipeline子协议协商、帧大小、空闲检测、背压水位都能逐项配。维度SSESpring 内置 WebSocketNetty WebSocket通信方向仅服务端到客户端双向双向连接规模万级偏吃力千级以内较稳万级起步可调心跳控制靠注释帧不可控需自己写IdleStateHandler 精确到秒子协议协商不支持部分支持支持 websocket subprotocol编排能力弱中强可插自定义 Codec典型场景公告、进度条后台管理通知IM、设备指令、大屏推送选 Netty 不是为了炫技而是当推送链路需要精确控制心跳间隔、单连接写缓冲上限、慢客户端降级策略时容器给的旋钮太少了。2.2 Netty 的 boss/worker 线程模型与 Spring 线程池的边界NioEventLoopGroup分成两组bossGroup 只负责 accept 新连接通常 1 个线程就够workerGroup 负责所有已建立连接的读写事件默认线程数是2 * CPU 核数。这里有一条最容易踩的红线不要在 ChannelHandler 里执行阻塞操作。数据库查询、Redis 调用、HTTP 远程调用都会卡住整个 EventLoop 上挂着的成百上千个连接。常见做法是 handler 里只做协议解析和状态维护业务逻辑通过ApplicationContext拿到 Spring 的 Service丢进业务线程池或者用CompletableFuture.supplyAsync异步执行。反过来说Channel.writeAndFlush()是线程安全的从任意业务线程调用都可以Netty 内部会把写任务提交到该 Channel 绑定的 EventLoop 上这一点是整套架构能成立的基础。2.3 会话路由表userId 到 Channel 的映射怎么写才不会内存泄漏路由表是 SpringBoot 业务层和 Netty 连接层之间唯一的桥。Component public class SessionManager { /** userId - Channel一个用户多端登录时可换成 userId : deviceId */ private final MapLong, Channel online new ConcurrentHashMap(); public void bind(Long userId, Channel channel) { Channel old online.put(userId, channel); // 同账号重复登录踢掉旧连接避免一条消息推两次 if (old ! null old ! channel old.isActive()) { old.writeAndFlush(new TextWebSocketFrame({\cmd\:\kick\})) .addListener(ChannelFutureListener.CLOSE); } // 把 userId 反向挂到 Channel 上channelInactive 时才知道该删哪条 channel.attr(USER_ID).set(userId); } public void unbind(Long userId, Channel channel) { // 用两参数版本防止把新连接误删 online.remove(userId, channel); } public Channel get(Long userId) { Channel ch online.get(userId); return (ch ! null ch.isActive()) ? ch : null; } public int size() { return online.size(); } }这段代码有三个细节值得单独说。第一online.remove(userId, channel)必须用两参数版本如果用户在同一瞬间重连旧连接的channelInactive回调可能晚于新连接的bind执行单参数 remove 会把新连接误删表现就是“刚连上就收不到消息”。第二USER_ID用AttributeKeyLong定义成静态常量比在 Handler 里再存一份 Map 更省内存也天然跟着 Channel 生命周期走。第三多实例部署时这张 Map 只覆盖本机连接用户连在 A 节点、消息从 B 节点发出就会丢需要在推送前先用 Redis 发布订阅广播一次各节点收到后再查本地路由表这一步不能省。3. 在 SpringBoot 里拉起 Netty WebSocket 服务端的最小可运行代码3.1 依赖与 application.yml 里的推送参数dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency dependency groupIdio.netty/groupId artifactIdnetty-all/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-data-redis/artifactId /dependencyNetty 的版本交给spring-boot-dependencies统一管理即可自己硬写版本号反而容易和 SpringBoot 里的其他 Netty 传递依赖打架出现NoSuchMethodError时先mvn dependency:tree | grep netty看一眼。push: websocket: port: 8888 path: /ws subprotocol: chat boss-threads: 1 worker-threads: 0 # 0 交给 Netty 默认即 2 * CPU max-frame-size: 65536 # 单个 WebSocket 帧上限 64KB http-max-content: 65536 # 握手阶段 HTTP 聚合上限 reader-idle-seconds: 90 # 超过 90 秒没收到客户端任何数据即判定假死reader-idle-seconds是整套推送里最该调准的参数设太短弱网下的手机端会被反复踢下线设太长服务端会长期挂着一批已经断网但没发 FIN 包的僵尸连接。经验值取业务心跳间隔的 2 到 3 倍。3.2 用 SmartLifecycle 托管 Netty 启动与优雅停机不要用PostConstruct启动 Netty。SpringBoot 的自动装配原理决定了容器刷新和 Bean 初始化有明确顺序而 Netty 需要的 Redis、数据库连接池如果还没就绪启动阶段就会抛异常。用SmartLifecycle更稳还能控制停机顺序。Component public class WebSocketServer implements SmartLifecycle { private EventLoopGroup boss; private EventLoopGroup worker; private Channel serverChannel; private volatile boolean running false; Override public void start() { boss new NioEventLoopGroup(props.getBossThreads()); worker new NioEventLoopGroup(props.getWorkerThreads()); ServerBootstrap b new ServerBootstrap(); b.group(boss, worker) .channel(NioServerSocketChannel.class) .option(ChannelOption.SO_BACKLOG, 1024) .childOption(ChannelOption.SO_KEEPALIVE, true) .childOption(ChannelOption.TCP_NODELAY, true) .childOption(ChannelOption.WRITE_BUFFER_WATER_MARK, new WriteBufferWaterMark(32 * 1024, 64 * 1024)) .childHandler(new WsChannelInitializer(props, sessionManager)); this.serverChannel b.bind(props.getPort()).syncUninterruptibly().channel(); this.running true; } Override public void stop() { if (serverChannel ! null) serverChannel.close(); // 先温和关闭超时再强制避免正在写的消息被截断 worker.shutdownGracefully(); boss.shutdownGracefully(); this.running false; } Override public boolean isRunning() { return running; } Override public int getPhase() { return Integer.MAX_VALUE; } // 最后启动最先停止 }getPhase()返回Integer.MAX_VALUE表示这个组件在容器启动的最后一刻被拉起、在停机的第一步被关闭正好符合“先断连接、再关资源”的诉求。syncUninterruptibly()保证绑定端口失败时直接抛异常让容器启动失败而不是静默留一个没监听的进程。3.3 握手阶段解析 tokenAuthHandshakeHandlerWebSocket 握手本质是一次 HTTP Upgrade 请求token 一般挂在 query string 或者Sec-WebSocket-Protocol里。public class AuthHandshakeHandler extends ChannelInboundHandlerAdapter { private final TokenService tokenService; Override public void channelRead(ChannelHandlerContext ctx, Object msg) { if (msg instanceof FullHttpRequest req) { Long userId tokenService.parse(uriParam(req.uri(), token)); if (userId null) { FullHttpResponse resp new DefaultFullHttpResponse( HttpVersion.HTTP_1_1, HttpResponseStatus.UNAUTHORIZED); ctx.writeAndFlush(resp).addListener(ChannelFutureListener.CLOSE); ReferenceCountUtil.release(msg); return; // 不往下传后面的 WebSocket 握手处理器不会看到这个请求 } ctx.channel().attr(USER_ID).set(userId); } ctx.fireChannelRead(msg); // 合法请求继续传给 WebSocketServerProtocolHandler } }ReferenceCountUtil.release(msg)必须写FullHttpRequest是引用计数的 ByteBuf 包装不释放就会内存泄漏而且泄漏报告只在-Dio.netty.leakDetection.levelparanoid下才会打出来。另外tokenService.parse如果设计成查库一定要在这里改成查本地缓存握手阶段还在 EventLoop 线程上。3.4 组装 pipeline 与 WebSocketFrameHandlerOverride protected void initChannel(SocketChannel ch) { ch.pipeline() .addLast(new HttpServerCodec()) .addLast(new HttpObjectAggregator(props.getHttpMaxContent())) .addLast(new AuthHandshakeHandler(tokenService)) .addLast(new IdleStateHandler(props.getReaderIdleSeconds(), 0, 0, TimeUnit.SECONDS)) .addLast(new WebSocketServerProtocolHandler( props.getPath(), props.getSubprotocol(), true, props.getMaxFrameSize())) .addLast(new WebSocketFrameHandler(sessionManager)); }WebSocketServerProtocolHandler的四个参数依次是路径、子协议、是否允许扩展allowExtensions压缩扩展需要它、单帧最大字节数。子协议填chat之后客户端握手时如果Sec-WebSocket-Protocol传了别的值握手会返回失败这是排查“Postman 能连、浏览器连不上”的第一现场。IdleStateHandler放在 WebSocket 处理器之前才能对所有阶段的读空闲统一生效。public class WebSocketFrameHandler extends SimpleChannelInboundHandlerWebSocketFrame { Override protected void channelRead0(ChannelHandlerContext ctx, WebSocketFrame frame) { if (frame instanceof TextWebSocketFrame text) { // 上行文本帧通常是心跳或业务指令 if (ping.equals(text.text())) { ctx.writeAndFlush(new TextWebSocketFrame(pong)); } } else if (frame instanceof CloseWebSocketFrame) { ctx.close(); } } Override public void channelInactive(ChannelHandlerContext ctx) { Long userId ctx.channel().attr(USER_ID).get(); if (userId ! null) sessionManager.unbind(userId, ctx.channel()); } Override public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { ctx.close(); } }SimpleChannelInboundHandler会自动释放消息不要重复 release。channelInactive是清理路由表唯一的可靠入口handlerRemoved在异常路径下不一定触发。4. 消息推送链路的打通从业务代码到 Channel 写出4.1 上下行消息的信封格式约定协议不用复杂但一定要有cmd和seq。前端靠cmd分发靠seq去重和排查乱序。{cmd:notice,seq:10241,ts:1735689600000,data:{title:工单待处理,id:8812}}上行的seq由客户端生成服务端回包时原样带上前端就能把请求和响应配成对不至于把心跳回包当成业务消息处理。4.2 MessagePushService业务代码怎么安全地 writeAndFlushService public class MessagePushService { private final SessionManager sessions; private final ObjectMapper mapper; public boolean pushToUser(Long userId, String cmd, Object data) { Channel ch sessions.get(userId); if (ch null) { return false; // 不在线交给调用方决定是否落离线表 } if (!ch.isWritable()) { log.warn(outbound buffer full, drop push. userId{}, userId); return false; // 背压保护宁可丢这条也不撑爆堆外内存 } String json toJson(new Envelope(cmd, data)); ch.writeAndFlush(new TextWebSocketFrame(json)).addListener(f - { if (!f.isSuccess()) { log.warn(push failed. userId{}, cause{}, userId, f.cause().getMessage()); } }); return true; } }三个判断缺一不可get()已经过滤了失活 ChannelisWritable()是背压闸门addListener是唯一能知道这条消息到底发出去没有的地方。业务侧拿到false时的常见做法是写一张离线消息表用户下次上线拉取最近 N 条而不是无脑重试。4.3 写缓冲水位与 isWritable 背压判断isWritable()的返回值由WriteBufferWaterMark决定水位线以下返回 true超过高水位返回 false回落到低水位之间才恢复。参数建议值作用低水位32KB低于此值 isWritable 恢复 true高水位64KB超过此值 isWritable 转 falseSO_BACKLOG1024已完成三次握手但未被 accept 的队列长度TCP_NODELAYtrue关闭 Nagle小消息即时发出max-frame-size64KB超过则解码抛 TooLongFrameExceptionreader-idle90s读空闲判定触发后发 Ping 或直接关闭水位线设太小正常突发推送就会误判为背压设太大慢客户端堆积的字节全压在堆外内存里直到 OOM 才被发现。64KB 高水位是多数文本推送场景的稳妥起点。4.4 心跳保活与 channelInactive 清理Override public void userEventTriggered(ChannelHandlerContext ctx, Object evt) { if (evt instanceof IdleStateEvent e e.state() IdleState.READER_IDLE) { ctx.writeAndFlush(new PingWebSocketFrame()); // 再给一次机会下一轮空闲直接关 if (missCount 2) ctx.close(); } else if (evt instanceof WebSocketServerProtocolHandler.HandshakeComplete) { missCount 0; sessionManager.bind(ctx.channel().attr(USER_ID).get(), ctx.channel()); } else { ctx.fireUserEventTriggered(evt); } }绑定的时机很关键必须等HandshakeComplete事件而不是channelActive或者认证通过那一刻。握手没完成就写 WebSocket 帧会被协议处理器直接丢弃。4.5 联调postman websocket 连接与 vue 侧重连Postman 新版本原生支持 WebSocket 请求直接填ws://127.0.0.1:8888/ws?tokenxxx然后在 Message 面板里手动发ping看是否能收到pong比写测试类快得多。如果是老版本只能装 JMeter 的 websocket sampler 插件来做压测和抓帧。Vue 侧用原生 WebSocket 加指数退避重连就够了let ws null, retry 0; function connect(token) { ws new WebSocket(ws://127.0.0.1:8888/ws?token${token}); ws.onopen () { retry 0; setInterval(() ws.send(ping), 30000); }; ws.onmessage (e) dispatch(JSON.parse(e.data)); ws.onclose () { // 指数退避上限 30 秒避免服务端重启时被瞬间打满 const delay Math.min(1000 * 2 ** retry, 30000); setTimeout(() connect(token), delay); }; }注意两点重连必须带上 token因为 token 可能已经刷新onmessage里要按cmd分发并做try/catch一条脏数据不能让整个消息循环挂掉。5. 粘包、半包与高并发下的参数调优5.1 netty粘包处理在 WebSocket 下的真实边界先说结论WebSocket 帧格式自带长度字段Netty 的WebSocketFrameDecoder已经帮你做了分帧重组不需要再往 pipeline 里加LengthFieldBasedFrameDecoder。真正会出现“粘包”感的地方只有两处。一是握手阶段。如果 token 或者 Cookie 特别大HttpObjectAggregator的maxContentLength不够就会抛TooLongFrameException客户端看到的是连接直接失败容易被误判成鉴权问题。二是你自己在文本帧里再套一层 JSON 数组批量发送时单帧超过max-frame-size会被拆成ContinuationWebSocketFrame这时必须自己判断isFinalFragment()并缓存中间分片否则每条消息都是半条 JSON表现为com.fasterxml.jackson解析报错。只有在完全不使用 WebSocket 协议、走裸 TCP 自定义协议时才需要LengthFieldBasedFrameDecoder来按长度切包。5.2 参数表与压测验证场景参数建议值万级连接worker-threadsCPU 核数默认 2 倍可能偏多推送延迟敏感TCP_NODELAYtrue大消息体max-frame-size256KB 或改用分片弱网移动端reader-idle120s 以上慢客户端保护WriteBufferWaterMark32KB / 64KB压测用 JMeter 的 websocket sampler把线程数当成并发连接数持续 10 分钟观察三件事SessionManager.size()是否和并发数一致少了说明丢连接多了说明僵尸连接没清、GC 是否出现长暂停、isWritable()为 false 的比例。如果第三条超过 1%说明消费者跟不上生产者要么加大 worker 线程要么在业务侧做推送限流。5.3 排错清单连不上、握手返回 401先看AuthHandshakeHandler有没有提前return却忘了ReferenceCountUtil.release再看 token 是否带了 URL 编码。连上后 30 秒断开对一下IdleStateHandler的读空闲和服务端心跳间隔客户端没回 Pong 就会被关。报TooLongFrameException同时调大HttpObjectAggregator.maxContentLength和WebSocketServerProtocolHandler.maxFrameSize只调一个没用。多实例下消息重复Redis 广播加上本机路由命中判断每个节点只推自己持有的那条 Channel。内存缓慢上涨开-Dio.netty.leakDetection.levelparanoid重点看AuthHandshakeHandler和自定义 Codec 的 release 路径。本文还有配套的精品资源点击获取
返回列表