ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

SpringBoot 整合 Netty WebSocket 实现高并发实时推送

2026/9/18 16:22:28 拓冰建站 浏览量
SpringBoot 整合 Netty WebSocket 实现高并发实时推送 简介面向 Java 后端开发者的消息推送实战示例适合已有一定 Spring 基础、希望把 WebSocket 长连接落地到真实项目的读者。内容以可运行的工程代码为主线先从依赖引入讲起再说明 NettyConfig 中 ChannelGroup 与用户-Channel 映射的搭建思路随后介绍 NettyServer 如何在独立线程中启动以避免阻塞主服务并借助 WebSocketHandler 处理连接建立、消息接收与断开时的清理与解绑。前端部分使用浏览器 WebSocket API 建立连接并渲染推送内容最终实现向全部在线用户广播、按用户 ID 定向推送两类典型场景关键步骤均配有注释便于对照调试和二次改造也有助于理解 Netty 的异步非阻塞模型与 Channel 生命周期。包内共 1 个 PDF 文件约 172KB属于图文结合的代码讲解型文档结构紧凑适合离线阅读与按需检索。当前已有 6276 人学习可作为梳理长连接管理方案的参考。1. 轮询和 SseEmitter 都试过之后为什么还是选 Netty 做推送做后台管理系统时运营要实时看到订单状态变化前端用 setInterval 每 3 秒拉一次接口十几个页面同时开着数据库 QPS 直接翻了几倍。换 SseEmitter 之后单机连接数一上来Tomcat 线程池就吃紧而且只支持服务端单向推前端想回传个 ACK 还得另开接口。于是把推送通道单独拆出来SpringBoot 负责业务和 HTTP 接口Netty 负责长连接和消息下发WebSocket 作为协议层承载。这套组合最大的价值在于推送粒度可控。ChannelGroup 拿来做全量广播ConcurrentHashMap 拿来做点对点定向推送两条路径互不干扰。适合需要给所有在线用户推公告又要给某个用户推私信的场景比如工单提醒、审批流节点通知、站内信。下面按依赖、容器、服务端启动、Handler 生命周期、接口验证的顺序拆一遍代码可以直接抄。2. netty-all 依赖与 ChannelGroup、ConcurrentHashMap 双索引容器2.1 依赖引入与版本取舍pom 里其实只需要两个依赖Netty 提供 WebSocket 服务端能力Hutool 只是拿来解析 JSON换成 Jackson 或 Fastjson 一样跑。!-- Netty 聚合包包含 transport、codec-http、handler 等 WebSocket 所需模块 -- dependency groupIdio.netty/groupId artifactIdnetty-all/artifactId version4.1.33.Final/version /dependency !-- 仅用到 JSONUtil可替换为 jackson-databind -- dependency groupIdcn.hutool/groupId artifactIdhutool-all/artifactId version5.2.3/version /dependencynetty-all是聚合包图省事可以直接引如果对包体积敏感实际只需要netty-transport、netty-codec-http、netty-handler三个。WebSocketServerProtocolHandler位于netty-codec-http下版本上 4.1.x 系列 API 基本稳定升到 4.1.100 也不用改代码。SpringBoot 版本本身不影响 Netty它俩是独立的依赖树spring-boot-starter-web 内嵌的 Tomcat 和 Netty 各占用自己的端口。2.2 两个容器各管一件事核心设计就一个类一个ChannelGroup管所有连接一个ConcurrentHashMap管用户和连接的映射关系。前者解决发给谁都不知道的广播问题后者解决发给某一个人的定位问题。public class NettyConfig { /** * 全局连接组管理所有在线 channel * GlobalEventExecutor.INSTANCE 是全局单例事件执行器不需要额外线程池 */ private static final ChannelGroup CHANNEL_GROUP new DefaultChannelGroup(GlobalEventExecutor.INSTANCE); /** * userId - Channel 映射用于定向推送 * 必须用 ConcurrentHashMap因为 Netty 的 IO 线程和业务线程会同时读写 */ private static final ConcurrentHashMapString, Channel USER_CHANNEL_MAP new ConcurrentHashMap(); private NettyConfig() { } public static ChannelGroup getChannelGroup() { return CHANNEL_GROUP; } public static ConcurrentHashMapString, Channel getUserChannelMap() { return USER_CHANNEL_MAP; } }DefaultChannelGroup内部同样是ConcurrentHashMap但它在add()时会给每个 channel 挂一个 closeFuture 监听器channel 一断开会自动从组里摘掉。这就是为什么广播路径不需要手动清理而USER_CHANNEL_MAP必须自己写remove——它没有这个自动机制。维度ChannelGroupConcurrentHashMap存储内容所有活跃 Channeluid 与 Channel 的映射写入时机handlerAddedchannelRead0 收到 uid清理方式自动监听 closeFuture手动 remove推送方式writeAndFlush 遍历广播get(uid) 后单发线程安全内部保证需 ConcurrentHashMap2.3 为什么不用 Map 直接装所有人有人会想一个MapString, Channel全搞定不行吗广播的时候遍历 value 就行。问题在于遍历过程中如果有 channel 关闭ConcurrentHashMap的弱一致性迭代虽然不抛异常但writeAndFlush会往已关闭的连接上写触发ClosedChannelException还得自己吞异常。ChannelGroup把批量写 自动剔除失效连接封装好了广播语义更干净。注意NettyConfig用的是静态字段 私有构造天然单例。别改成 Spring Bean 再注入到 Handler 里Handler 本身已经是单例两边生命周期容易打架。3. NettyServer 的 pipeline 编排与 Spring 生命周期托管3.1 boss 与 worker 的线程职责Netty 的 Reactor 模型在这里体现为两个EventLoopGroup。boss 只管accept新连接拿到SocketChannel后注册到 worker 上worker 负责后续所有读写事件。boss 线程数设 1 就够因为 accept 操作极快worker 默认是2 * CPU 核心数通常在 4 核机器上跑几百个长连接毫无压力。bootstrap.bind().sync()之后的channelFuture.channel().closeFuture().sync()是阻塞调用它会一直挂在那里等待服务端通道关闭。这两行合起来意味着start()方法不会返回在主线程里直接调就把 SpringBoot 启动流程卡死了RestController全都注册不上。3.2 pipeline 里每个 Handler 的顺序不能乱WebSocket 握手本质上是一次 HTTP 请求带着Upgrade: websocket头。所以链路必须先把 HTTP 协议解出来再由协议处理器完成升级。顺序如下ChannelHandler.Sharable Component public class WebSocketHandler extends SimpleChannelInboundHandlerTextWebSocketFrame { }Component public class NettyServer { private static final String WEBSOCKET_PROTOCOL WebSocket; Value(${webSocket.netty.port:58080}) private int port; Value(${webSocket.netty.path:/webSocket}) private String webSocketPath; Autowired private WebSocketHandler webSocketHandler; private EventLoopGroup bossGroup; private EventLoopGroup workGroup; private void start() throws InterruptedException { bossGroup new NioEventLoopGroup(1); workGroup new NioEventLoopGroup(); ServerBootstrap bootstrap new ServerBootstrap(); bootstrap.group(bossGroup, workGroup) .channel(NioServerSocketChannel.class) .localAddress(new InetSocketAddress(port)) .childHandler(new ChannelInitializerSocketChannel() { Override protected void initChannel(SocketChannel ch) { // 1. HTTP 编解码WebSocket 握手走的是 HTTP 请求 ch.pipeline().addLast(new HttpServerCodec()); // 2. 分块写处理大报文避免一次性堆内存 ch.pipeline().addLast(new ChunkedWriteHandler()); // 3. 聚合 HTTP 分段8192 是单次聚合的最大字节数 ch.pipeline().addLast(new HttpObjectAggregator(8192)); // 4. 协议升级path 必须与前端 ws URL 路径完全一致 ch.pipeline().addLast(new WebSocketServerProtocolHandler( webSocketPath, WEBSOCKET_PROTOCOL, true, 65536 * 10)); // 5. 业务 handler 放最后前面都是协议层 ch.pipeline().addLast(webSocketHandler); } }); ChannelFuture future bootstrap.bind().sync(); future.channel().closeFuture().sync(); } }HttpObjectAggregator(8192)的参数是内容长度上限超了会返回 413。WebSocketServerProtocolHandler的第四个参数65536 * 10是单帧最大字节数默认只有 64KB推送长文本或 base64 图片时很容易踩到TooLongFrameException这里放大到 640KB。第三个参数allowExtensions传 true允许客户端协商扩展。原示例里还挂了ObjectEncoder这是 Java 序列化编码器作用是把 POJO 编码成 ByteBuf。但推送用的是TextWebSocketFrame走的是文本帧通道这个 handler 实际不会命中属于冗余配置删掉不影响功能。3.3 用 Spring 生命周期接管启停PostConstruct在 Bean 初始化后执行在这里另起线程跑start()主线程就能继续完成 SpringBoot 的其余装配。PostConstruct public void init() { Thread serverThread new Thread(() - { try { start(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } }, netty-websocket-server); // 设为守护线程主进程退出时不残留 serverThread.setDaemon(true); serverThread.start(); } PreDestroy public void destroy() throws InterruptedException { if (bossGroup ! null) { bossGroup.shutdownGracefully().sync(); } if (workGroup ! null) { workGroup.shutdownGracefully().sync(); } }PreDestroy在容器关闭时触发shutdownGracefully()会先拒绝新任务再把队列里的任务跑完默认静默期 2 秒、超时 15 秒。加了.sync()是为了等线程池真正退干净再往下走避免热部署时端口被旧进程占着。提示new Thread不设名字的话jstack 里只能看到Thread-N线上排查连接堆积时很难定位。给线程起名是成本最低的可观测性投入。端口和路径通过Value外置启动参数--webSocket.netty.port58081就能覆盖多环境部署不用改代码。4. WebSocketHandler 的连接生命周期与用户身份绑定4.1 Sharable 与泛型选择Handler 被声明为 Spring 单例而每个连接都会把它加进自己的 pipeline。如果没标Sharable第二个客户端连上来时 Netty 会直接抛is not a Sharable handler, so cant be added or removed multiple times连接建立失败。泛型选TextWebSocketFrame而不是WebSocketFrame是因为推送内容都是 JSON 字符串。文本帧会被自动转型后传进channelRead0不用再手动判断帧类型。4.2 三个回调构成完整生命周期Component ChannelHandler.Sharable public class WebSocketHandler extends SimpleChannelInboundHandlerTextWebSocketFrame { private static final Logger log LoggerFactory.getLogger(WebSocketHandler.class); private static final AttributeKeyString USER_ID_KEY AttributeKey.valueOf(userId); /** 连接建立第一个被调用 */ Override public void handlerAdded(ChannelHandlerContext ctx) { log.info(连接建立 channelId{}, ctx.channel().id().asLongText()); NettyConfig.getChannelGroup().add(ctx.channel()); } /** 收到客户端文本帧 */ Override protected void channelRead0(ChannelHandlerContext ctx, TextWebSocketFrame msg) { JSONObject json JSONUtil.parseObj(msg.text()); String uid json.getStr(uid); if (uid null || uid.isEmpty()) { ctx.channel().writeAndFlush(new TextWebSocketFrame(uid 不能为空)); return; } // 建立 uid 与 channel 的映射 NettyConfig.getUserChannelMap().put(uid, ctx.channel()); // 写入 channel 属性断开时靠它反查 uid ctx.channel().attr(USER_ID_KEY).set(uid); ctx.channel().writeAndFlush(new TextWebSocketFrame(服务器连接成功)); } /** 连接断开 */ Override public void handlerRemoved(ChannelHandlerContext ctx) { NettyConfig.getChannelGroup().remove(ctx.channel()); removeUserId(ctx); log.info(连接断开 channelId{}, ctx.channel().id().asLongText()); } Override public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { log.warn(连接异常{}, cause.getMessage()); NettyConfig.getChannelGroup().remove(ctx.channel()); removeUserId(ctx); ctx.close(); } private void removeUserId(ChannelHandlerContext ctx) { String userId ctx.channel().attr(USER_ID_KEY).get(); if (userId ! null) { NettyConfig.getUserChannelMap().remove(userId); } } }handlerAdded在 channel 注册进 pipeline 时触发此时连接已就绪适合做注册动作。channelRead0是唯一的入站数据处理入口SimpleChannelInboundHandler会自动release()消息对象不需要手动释放引用计数。handlerRemoved在 channel 从 pipeline 移除时触发正常关闭和异常关闭都会走到。4.3 AttributeKey 的复用与读取AttributeKey.valueOf(userId)内部维护了一个常量池同一个字符串多次调用返回的是同一个实例所以放在静态常量里复用即可不要用newInstance()重名会直接抛异常。属性存在 channel 上而不是 Handler 上这样即使多个连接共用一个 Handler 实例每个连接的 uid 也是隔离的。注意原示例用的是setIfAbsent(uid)。如果同一个连接先后发了两次不同 uid属性会保留第一次的值而userChannelMap已经被第二次覆盖。结果是断开时removeUserId删掉的是旧 uid新 uid 的记录残留在 map 里形成永远推不到但也不报错的僵尸条目。直接用set()覆盖更安全。4.4 同一账号多端登录的覆盖问题userChannelMap.put(uid, channel)是覆盖写。手机端和网页端用同一个 uid 登录时后连的会把前面的挤掉。更麻烦的是当前一个连接断开时handlerRemoved会执行remove(uid)把新的那条记录也删了导致在线用户收不到推送。常见做法是把 value 换成Channel的集合或者加一个判断只在当前 channel 与 map 中记录一致时才删除private void removeUserId(ChannelHandlerContext ctx) { String userId ctx.channel().attr(USER_ID_KEY).get(); if (userId null) { return; } // 只移除属于自己的那条记录避免顶掉新连接 NettyConfig.getUserChannelMap().remove(userId, ctx.channel()); }ConcurrentHashMap.remove(key, value)是原子操作只有 value 相等才删除。这一步改动很小但能挡掉多端登录场景下的绝大部分推送丢失工单。5. 推送接口落地、Postman 端到端验证与排错清单推送逻辑封装成一个 Service对外暴露两个方法Controller 层只做参数接收。Service public class PushServiceImpl implements PushService { Override public void pushMsgToOne(String userId, String msg) { Channel channel NettyConfig.getUserChannelMap().get(userId); if (channel null || !channel.isActive()) { // 用户不在线落库或走离线消息队列 throw new IllegalStateException(用户 userId 不在线); } channel.writeAndFlush(new TextWebSocketFrame(msg)); } Override public void pushMsgToAll(String msg) { TextWebSocketFrame frame new TextWebSocketFrame(msg); // ChannelGroup 内部自动跳过已关闭的连接 NettyConfig.getChannelGroup().writeAndFlush(frame); } }channel.isActive()这个判断不能省。map 里的记录可能因为handlerRemoved没来得及执行而短暂残留直接往失效连接写会抛ClosedChannelException。广播路径之所以不用判是因为ChannelGroup在写之前会检查每个 channel 的状态。前端页面拿到消息后在onmessage里追加到 textarea 即可注意onopen中先把 uid 发过去否则服务端拿不到用户标识。验证分两步。先在浏览器打开 HTML 页面确认控制台收到服务器连接成功服务端日志打印出{uid:123456}。再用 Postman 发 POST 请求/push/pushAll?msg全体通知和/push/pushOne?userId123456msg定向通知两种粒度各测一次。Postman 自带的 WebSocket 客户端也能用来手工连接ws://ip:58080/webSocket方便观察原始帧内容。现象大概率原因排查点握手返回 404路径不匹配前端 ws URL 路径与webSocketPath是否一致连接立即 1006 关闭pipeline 顺序错或帧超限HttpObjectAggregator是否在协议处理器之前广播正常、定向失败uid 未上报或已被覆盖userChannelMap中是否存在该 key推送后前端无反应channel 已失效检查isActive()与服务端日志的 handlerRemoved内存持续增长map 未清理removeUserId是否被exceptionCaught覆盖到生产环境还需要补两件事一是心跳WebSocketServerProtocolHandler已内置 ping/pong 处理客户端侧加个 30 秒定时 ping 就能让 Nginx 的 60 秒空闲断连不再触发二是推送成功率监控在writeAndFlush返回的ChannelFuture上加addListener统计失败次数比等用户投诉要早得多。本文还有配套的精品资源点击获取