资讯详情

SpringBoot+Netty WebSocket长连接推送全链路解析

📅 2026/9/18 7:27:15 | 华诺云谱 👁 阅读
SpringBoot+Netty WebSocket长连接推送全链路解析
简介这份资源面向具备一定Java与网络编程基础的开发者围绕SpringBoot、WebSocket与Netty三大组件给出可运行的消息推送示例代码解决服务端主动向客户端实时推送消息、并区分全员广播与指定用户定向推送的问题适用于聊天、通知提醒、行情报价等实时交互场景。压缩包内共1个PDF文件约172KB以图文结合的代码讲解方式呈现从依赖引入、连接管理到处理器实现逐步展开。内容涵盖在NettyConfig中借助ChannelGroup统一管理所有channel、用ConcurrentHashMap维护用户ID与channel的映射关系创建包含bossGroup与workerGroup的NettyServer并在独立线程中启动以避免阻塞SpringBoot主服务实现继承SimpleChannelInboundHandler的WebSocketHandler处理连接建立、消息接收与关闭事件前端通过JavaScript的WebSocket API建立连接并监听message事件展示推送内容。读者可据此梳理清消息推送的整体链路掌握广播与定向发送两种调用方式并理解Netty异步非阻塞模型在实时通信中的作用已有6276人学习。1. 为什么要在 SpringBoot 里再拉一条 Netty 的 WebSocket 通道很多人在 IDEA 创建 SpringBoot 项目后第一反应是引 spring-boot-starter-websocket用ServerEndpoint把推送做掉。几百个连接时这条路毫无问题连接数一上来就会撞到天花板会话状态散落在各节点本地广播一次要扫全量连接长连接还和业务 HTTP 线程池抢资源。更常见的做法是让 SpringBoot 管业务和 HTTP 接口另起一个 Netty 端口专门跑 WebSocket 长连接两边用同一个 JVM 内的 Service 调用或跨节点的 Redis 打通。它解决的是「服务端主动、低延迟、大连接数」这一类推送诉求站内信、订单状态流转、工单提醒、监控告警都算。适合已经用 SpringBoot 做后端、又不想为推送单独引入一套消息中间件的团队。2. Netty 服务端骨架从 ServerBootstrap 到 WebSocketServerProtocolHandlerSpringBoot 自带的内嵌容器本质是 Servlet 模型WebSocket 连接最终会被包装成 servlet 会话扩容时会话不共享。Netty 走的是 NIO 事件循环一个 EventLoop 线程能扛住成千上万个 Channel写数据和读数据都在同一线程里串行执行天然避开了锁竞争。这一章把一条能跑的 Netty WebSocket 服务端从依赖到优雅停机拆开讲清楚代码可以直接落到项目里再改。2.1 依赖与端口规划netty-all不在 Spring Boot 的 dependencyManagement 里版本必须自己锁死否则升 Spring Boot 时可能被动换版本。!-- pom.xml -- properties netty.version4.1.x/netty.version /properties dependencies dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency !-- NettyWebSocket 服务端 事件循环 -- dependency groupIdio.netty/groupId artifactIdnetty-all/artifactId version${netty.version}/version /dependency !-- 推送前先落离线消息、多实例广播时用 -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-data-redis/artifactId /dependency /dependencies端口不要和 8080 抢。常见做法是业务 HTTP 留在 8080Netty 单独监听一个端口比如 9000。端口协议承载内容是否对外暴露8080HTTP业务 REST 接口、推送触发接口经网关暴露9000TCP / WebSocket长连接只处理/ws路径经网关升级暴露9001HTTPNetty 健康检查、连接数指标仅内网application.yml里把端口做成配置容器化部署时好覆盖push: ws: port: 9000 path: /ws reader-idle-seconds: 60参数说明port是 Netty 监听端口改它不用改代码path是握手路径网关转发规则要和它对齐reader-idle-seconds决定多久没收到任何数据就判定连接已死前端心跳间隔必须比它小。2.2 ChannelPipeline 的处理器顺序怎么排Pipeline 是责任链顺序错了握手直接失败。正确顺序是空闲检测 → HTTP 编解码 → HTTP 聚合 → 鉴权 → WebSocket 协议处理 → 业务处理。Component public class WebSocketServer implements SmartLifecycle { Value(${push.ws.port:9000}) private int port; private EventLoopGroup bossGroup; private EventLoopGroup workerGroup; private Channel serverChannel; private volatile boolean running false; Override public void start() { bossGroup new NioEventLoopGroup(1); // 只管 accept1 个线程够用 workerGroup new NioEventLoopGroup(); // 默认 CPU 核数 * 2 ServerBootstrap bootstrap new ServerBootstrap(); bootstrap.group(bossGroup, workerGroup) .channel(NioServerSocketChannel.class) .option(ChannelOption.SO_BACKLOG, 1024) // 半连接队列 .childOption(ChannelOption.SO_KEEPALIVE, true) .childOption(ChannelOption.TCP_NODELAY, true) // 小包立即发别攒 .childHandler(new ChannelInitializerSocketChannel() { Override protected void initChannel(SocketChannel ch) { ch.pipeline() // 1. 读空闲检测超时触发 userEventTriggered .addLast(new IdleStateHandler(60, 0, 0, TimeUnit.SECONDS)) // 2. WebSocket 握手本质是一个 HTTP 升级请求 .addLast(new HttpServerCodec()) .addLast(new HttpObjectAggregator(64 * 1024)) // 3. 鉴权必须排在协议处理器前面 .addLast(new AuthHandler()) // 4. 协议处理器checkStartsWith 让 /ws?tokenxxx 也能匹配 .addLast(new WebSocketServerProtocolHandler( WebSocketServerProtocolConfig.newBuilder() .websocketPath(/ws) .checkStartsWith(true) .maxFrameSize(64 * 1024) .build())) // 5. 业务帧处理 .addLast(new PushMessageHandler()); } }); try { serverChannel bootstrap.bind(port).sync().channel(); running true; } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new IllegalStateException(netty 启动失败, e); } } }这里有两个容易踩的点。第一checkStartsWith(true)必须开。Netty 判断请求路径时默认用uri.equals(websocketPath)而浏览器握手发过来的是/ws?tokenabcquery 一挂上去就匹配不上握手直接被拒。第二HttpObjectAggregator不能省握手请求头加 upgrade 后是完整 HTTP 报文不聚合会被拆成多段协议处理器拿不到完整请求。2.3 握手鉴权token 放在 URL 还是 Sec-WebSocket-Protocol浏览器的WebSocket构造函数不支持自定义请求头所以 Cookie 和 URL query 是最常见的两种携带凭证方式。URL query 简单直接但 token 会进网关访问日志记得用短时效 token。Sec-WebSocket-Protocol更干净前端第二个参数传数组即可代价是服务端要显式配subprotocols不匹配会协商失败。public class AuthHandler extends ChannelInboundHandlerAdapter { Override public void channelRead(ChannelHandlerContext ctx, Object msg) { if (msg instanceof FullHttpRequest) { FullHttpRequest req (FullHttpRequest) msg; String token new QueryStringDecoder(req.uri()) .parameters().getOrDefault(token, List.of()).get(0); Long userId TokenService.parse(token); // 校验失败返回 null if (userId null) { FullHttpResponse resp new DefaultFullHttpResponse( HttpVersion.HTTP_1_1, HttpResponseStatus.UNAUTHORIZED, Unpooled.copiedBuffer(token invalid, StandardCharsets.UTF_8)); resp.headers().set(HttpHeaderNames.CONTENT_TYPE, text/plain; charsetutf-8); ctx.writeAndFlush(resp).addListener(ChannelFutureListener.CLOSE); ReferenceCountUtil.release(msg); // 自己消费了必须释放 return; } ctx.channel().attr(AttrKeys.USER_ID).set(userId); // 挂到 Channel 上 } ctx.fireChannelRead(msg); // 交棒不要在这里 release } }逻辑说明FullHttpRequest只在这一层出现放行后所有权归下一个 handler所以只有在拒绝分支里才release放行分支直接fireChannelRead。把userId写成AttributeKeyLong挂在 Channel 上后续 handler 用ctx.channel().attr(AttrKeys.USER_ID).get()取值比再解析一次 URI 便宜得多。2.4 交给 Spring 容器管生命周期Netty 的线程组必须随应用启停自己new Thread起服务会漏掉停机钩子。实现SmartLifecycle是最省事的做法getPhase返回Integer.MAX_VALUE表示最后启动、最先停止。Override public void stop() { if (serverChannel ! null) { serverChannel.close(); // 停止接收新连接 } if (workerGroup ! null) { workerGroup.shutdownGracefully(2, 5, TimeUnit.SECONDS); // 给在途消息留 2 秒 } if (bossGroup ! null) { bossGroup.shutdownGracefully(2, 5, TimeUnit.SECONDS); } running false; } Override public boolean isRunning() { return running; } Override public int getPhase() { return Integer.MAX_VALUE; }shutdownGracefully(quietPeriod, timeout, unit)的前两个参数是静默期和超时时间。静默期内如果没有新任务才真正退出给正在写出的推送帧留出冲刷时间。K8s 滚动发布时这个 2 秒能挡掉一批「连接被重置」的报错。3. 连接注册表与推送链路ChannelGroup、粘包处理与心跳服务端跑起来只完成一半真正决定推送系统能不能用的是连接注册表设计、消息分帧处理和断线判定这三件事。这一章给出可以直接抄的注册表实现、推送接口写法以及 netty 粘包处理在 WebSocket 场景下到底还要不要做。3.1 用 ConcurrentHashMap 建 userId 到 Channel 的映射点对点推送要求 O(1) 找到连接广播又要求能遍历全部连接所以两张表各管一段ConcurrentHashMap管定向查找ChannelGroup管广播。Component public class ConnectionManager { private static final Logger log LoggerFactory.getLogger(ConnectionManager.class); /** userId - Channel用于点对点推送 */ private final MapLong, Channel online new ConcurrentHashMap(); /** 全部活跃连接用于广播ChannelGroup 会自动剔除已关闭的 Channel */ private final ChannelGroup all new DefaultChannelGroup(GlobalEventExecutor.INSTANCE); public void register(Long userId, Channel channel) { Channel old online.put(userId, channel); if (old ! null old ! channel old.isActive()) { old.close(); // 同账号多端登录策略踢掉旧的 } all.add(channel); // 连接关闭时自动清理避免 Map 里留下僵尸 Channel 造成内存泄漏 channel.closeFuture().addListener(f - online.remove(userId, channel)); log.info(ws online userId{}, total{}, userId, all.size()); } public boolean sendToUser(Long userId, String json) { Channel ch online.get(userId); if (ch null || !ch.isActive()) { return false; } ch.writeAndFlush(new TextWebSocketFrame(json)); return true; } public void broadcast(String json) { all.writeAndFlush(new TextWebSocketFrame(json)); } public int onlineCount() { return all.size(); } }关键点是closeFuture().addListener。只靠channelInactive回调也能清理但连接在注册之前就断开的极端时序下会漏加一道兜底更稳。ChannelGroup.writeAndFlush内部会跳过已关闭的连接不用自己判活。3.2 netty粘包处理在 WebSocket 场景下到底还要不要做这是问得最多的问题。答案分两层TCP 裸连接上跑自定义协议粘包必须自己处理靠长度字段拆// 裸 TCP 场景4 字节长度头 业务体 p.addLast(new LengthFieldBasedFrameDecoder( 1024 * 1024, // maxFrameLength单帧上限防止恶意大包打爆内存 0, // lengthFieldOffset长度字段从第 0 字节开始 4, // lengthFieldLength长度字段占 4 字节 0, // lengthAdjustment长度值不包含头本身所以偏移为 0 4)); // initialBytesToStrip把 4 字节长度头剥掉再交给业务 p.addLast(new LengthFieldPrepender(4)); // 编码端对称补头而 WebSocket 本身是消息边界协议WebSocketServerProtocolHandler在maxFrameSize 0时会自动挂上WebSocketFrameAggregator把分片的 continuation frame 聚合成完整帧。所以标准 WebSocket 文本推送不需要额外做粘包处理写了反而多余。真正需要补的是第三种情况业务自定义二进制协议塞进BinaryWebSocketFrame。这时 WebSocket 只保证「一个 frame 是一整块」frame 内部还能再拆协议头// 在 WebSocket 帧之后再对 payload 做一次长度拆包 public static byte[] unpack(ByteBuf payload) { int length payload.readInt(); // 前 4 字节是业务长度 if (length 0 || length payload.readableBytes()) { throw new IllegalArgumentException(非法长度: length); } byte[] body new byte[length]; payload.readBytes(body); return body; }参数说明maxFrameSize建议设成业务最大消息体的 2 到 4 倍设太小会把正常的大 JSON 消息切断设太大等于给内存攻击留了口子。一般文本推送 64KB 足够。3.3 从 SpringBoot 的 HTTP 接口推到 Netty 的 Channel推送入口留在 Spring MVC 里业务代码只用调 Service不感知 Netty 的存在。RestController RequestMapping(/api/push) public class PushController { private final ConnectionManager connectionManager; private final OfflineMessageService offlineService; public PushController(ConnectionManager cm, OfflineMessageService os) { this.connectionManager cm; this.offlineService os; } PostMapping(/user/{userId}) public MapString, Object pushToUser(PathVariable Long userId, RequestBody PushMsg msg) { String payload JSON.toJSONString(msg); boolean delivered connectionManager.sendToUser(userId, payload); if (!delivered) { offlineService.save(userId, msg); // 不在线就落库上线后拉取 } return Map.of(delivered, delivered, msgId, msg.getId()); } }逻辑说明先在本地连接表里找找到就写帧写不出去或找不到就转离线。注意writeAndFlush是异步的返回 true 只代表消息进了发送队列不代表客户端收到了 ACK。真正确认送达要靠客户端回执这一点在第 5 章展开。3.4 心跳、IdleStateHandler 与前端断线重连IdleStateHandler只负责「多久没动静就触发事件」怎么处理由业务决定。参数建议值作用调整影响readerIdleTime60s读空闲判定调小发现死连接更快弱网误杀概率上升前端心跳间隔25s主动发 ping 帧必须小于 readerIdleTimewriterIdleTime0一般不启用空闲时自动发 ping 会干扰业务计数重连退避上限30s前端重连间隔封顶不封顶会在服务端故障时形成重连风暴Override public void userEventTriggered(ChannelHandlerContext ctx, Object evt) { if (evt instanceof IdleStateEvent) { ctx.close(); // 超时直接断开等前端重连不要试图发心跳探活 } else if (evt instanceof WebSocketServerProtocolHandler.HandshakeComplete) { Long userId ctx.channel().attr(AttrKeys.USER_ID).get(); connectionManager.register(userId, ctx.channel()); } } Override public void channelInactive(ChannelHandlerContext ctx) { Long userId ctx.channel().attr(AttrKeys.USER_ID).get(); log.info(ws offline userId{}, userId); }前端重连用指数退避加随机抖动别用固定 1 秒let retry 1000; function connect() { const ws new WebSocket(wss://host/ws?token${token}); ws.onopen () { retry 1000; }; ws.onclose () { // 抖动避免所有客户端同时重连 setTimeout(connect, retry Math.random() * 500); retry Math.min(retry * 2, 30000); }; ws.onmessage (e) { /* 解析业务消息回 ACK */ }; }4. 联调与排错postman websocket连接、握手 400 与多实例广播代码写完只是开始握手失败、消息收不到、多实例只推一台机器这三类问题占了排查时间的大头。这一章给出从单机联调到多实例部署的完整排查路径。4.1 用 Postman 和浏览器控制台把握手跑通Postman 新建请求时选 WebSocket 类型地址填ws://127.0.0.1:9000/ws?token测试token点 Connect。能连上说明握手链路通了再点发送发一条文本帧验证上行。命令行下用curl也能看握手头curl -i -N \ -H Connection: Upgrade \ -H Upgrade: websocket \ -H Sec-WebSocket-Version: 13 \ -H Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ \ http://127.0.0.1:9000/ws?tokentest正常返回HTTP/1.1 101 Switching Protocols响应头里带Upgrade: websocket和Sec-WebSocket-Accept。返回 400 一般是路径或者头不对返回 401 就是自己的鉴权拦下来的。4.2 握手失败的排查路径与常见异常栈现象最可能的原因定位方式400 Bad Request请求 URL 没匹配上 websocketPath检查是否开了checkStartsWith(true)401 Unauthorizedtoken 解析失败或已过期在 AuthHandler 里打断点看 URI403 Forbidden网关或鉴权中间件先拦了看网关访问日志比对转发规则101 成功但立刻断开业务 handler 里抛了异常看exceptionCaught有没有吞掉堆栈连上收不到消息注册发生在 HandshakeComplete 之前失败在register里打日志看 userId 是否为 null最常见的异常栈是java.lang.IllegalArgumentException: unsupported message type: class io.netty.handler.codec.http.DefaultFullHttpRequest。原因通常是业务 handler 声明的泛型是TextWebSocketFrame但握手阶段的 HTTP 请求也流经了同一条链。检查方式是确认WebSocketServerProtocolHandler是否真的拦截并消费了握手请求——路径没匹配上时它不消费请求就漏到后面去了这正是checkStartsWith那个坑的连锁反应。4.3 Nacos 多实例下推送只到一台机器怎么办Nacos 做服务发现时长连接是直接打到具体实例上的网关转发到 A 实例用户在 A 上有连接业务服务 B 实例调推送接口B 的本地连接表里没有这个 userId消息就丢了。Netty 本身不解决跨实例广播它只管单机连接。思路是把「本地能发就发发不了广播」做成默认行为Service public class ClusterPushService { private static final String TOPIC ws:push; private final ConnectionManager connectionManager; private final StringRedisTemplate redis; public void push(Long userId, PushMsg msg) { // 1. 先试本地命中率通常在 1/N ~ 100%取决于网关策略 if (connectionManager.sendToUser(userId, JSON.toJSONString(msg))) { return; } // 2. 本地没有广播给所有实例谁有这个连接谁来发 redis.convertAndSend(TOPIC, JSON.toJSONString( Map.of(userId, userId, msg, msg))); } }订阅端每个实例都收到各自查自己的连接表命中才发送Component public class PushMessageListener implements MessageListener { Override public void onMessage(Message message, byte[] pattern) { String body new String(message.getBody(), StandardCharsets.UTF_8); ClusterMsg cm JSON.parseObject(body, ClusterMsg.class); // 没有本地连接就静默返回天然幂等 connectionManager.sendToUser(cm.getUserId(), JSON.toJSONString(cm.getMsg())); } }要注意广播是 N 倍流量实例数上百时 Redis pub/sub 会成为瓶颈这时候改用按 userId 哈希的定向队列更合适。压测时观察all.size()和 Redis 订阅端 QPS两者能对上说明链路正常。5. 消息可靠性补强离线消息、推送幂等与 SSE 的选型边界单机连接数上来之后掉消息和重复推是两个必然出现的问题前者靠离线补偿后者靠幂等键。5.1 离线消息与 ACK 重投推送成功的判断不能只看writeAndFlush的返回值。让客户端收到消息后回一条 ACK 帧服务端把消息放进待确认表超时未确认就重投最多重投两次超过就落离线库。CREATE TABLE ws_offline_msg ( id BIGINT PRIMARY KEY AUTO_INCREMENT, user_id BIGINT NOT NULL, msg_id VARCHAR(64) NOT NULL, payload TEXT NOT NULL, status TINYINT NOT NULL DEFAULT 0, -- 0 待投递 1 已投递 2 已确认 retry_count INT NOT NULL DEFAULT 0, create_time DATETIME NOT NULL, UNIQUE KEY uk_msg (msg_id) -- 幂等键防重复投递 );uk_msg这个唯一索引就是幂等的兜底。客户端重连后先拉一次未确认消息拉到之后按msg_id去重再渲染服务端即使重复投递也不会造成界面重复。5.2 什么时候该用 SSE 而不是 WebSocket如果推送方向只有服务端到客户端业务不需要客户端上行SSE 是更省事的选择它跑在 HTTP 上网关、鉴权、日志全部复用现有链路浏览器EventSource自带断线重连和Last-Event-ID断点续传服务端一个SseEmitter就能推。代价是单向、纯文本、HTTP/1.1 下同源连接数有上限HTTP/2 下这个限制基本消失。判断标准很直接判断项选 WebSocket选 SSE客户端需要频繁上行是否消息体是二进制是否需要复用现有 HTTP 鉴权链路否是希望浏览器自动重连、续传否是单机连接数目标十万级万级压测验证时先用最笨的办法写个 Java 客户端循环建一万个连接观察all.size()、文件描述符数量和 Full GC 频率。如果 fd 增长和连接数不成正比说明有 Channel 没被正确释放重点查closeFuture回调和ChannelGroup是不是漏了清理。本文还有配套的精品资源点击获取
📝

华诺云谱内容团队

资深建站顾问 · 行业研究员

10年+企业数字化服务经验,专注智能建站、SEO优化与品牌营销,持续输出建站技巧、行业洞察与营销干货,已帮助5000+企业实现数字化增长。

你可能需要的服务

订阅华诺云谱资讯周报

每周一封,精选建站技巧、SEO与营销干货,直达邮箱。已有 8,000+ 企业主订阅,助你少走弯路。