简介:这份PDF格式的示例代码资源,完整演示了在SpringBoot项目中利用Netty作为后台服务端、前端通过WebSocket建立长连接的消息推送实现方案。资源面向具备一定Java基础、希望快速上手实时通信开发的读者,重点解决了服务端向全体用户广播以及按用户ID定向推送这两类常见需求。包内为1个PDF文件,体积约172KB,内容覆盖依赖引入、NettyConfig通道组与用户映射表设计、NettyServer双EventLoopGroup线程模型、WebSocketHandler连接建立与关闭处理,以及前端页面收发消息逻辑等模块,关键代码均配有注释。通过阅读可以掌握ChannelGroup广播机制、ConcurrentHashMap维护用户与连接对应关系、独立线程启动Netty避免阻塞主服务等实践要点,同时理解WebSocket协议在客户端与服务端之间的通信规范,便于快速迁移到聊天室、通知推送、设备状态监控等实时交互场景。目前已有6276人学习,对希望借鉴成熟实现、减少重复踩坑的开发者有直接参考价值。
1. WebSocket推送方案那么多,为什么偏偏绕不开Netty
做消息推送的第一反应往往是“SpringBoot自带的WebSocket不是开箱即用吗”,可真到了线上,要么是推送延迟从毫秒级变成秒级,要么是连接数一过几千就频繁掉线,要么是后端想主动往某个会话塞消息时找不到那个Channel。这个时候再看SpringBoot+WebSocket+Netty的示例代码,核心并不是“谁比谁高级”,而是把连接管理、线程模型、协议编解码从黑盒里捞出来,让推送链路的每一个环节都可控。这篇文章面对的读者是那些已经能用SpringBoot写CRUD、但对长连接心有顾虑的开发者:我会先把WebSocket和Netty的职责边界讲清楚,再给一套能直接跑的最小工程,最后把连接失效、线程阻塞、多实例推送这些常见坑一个个拆开。目录看完,你会知道这套代码值不值得放进自己的项目。
2. 先立住理论:WebSocket协议、SpringBoot内嵌容器与Netty的边界
2.1 WebSocket从握手到帧:推送实时消息到底在推什么
WebSocket本质上是一条建立在TCP之上的、由HTTP Upgrade“临时转行”的通信管道。客户端先发一个带有Upgrade: websocket和Sec-WebSocket-Key的HTTP请求,服务端返回101 Switching Protocols,之后双方就不再走HTTP语义,而是直接收发WebSocket帧。每一个帧由FIN、opcode、mask、payload length和payload组成。服务端发给客户端的帧不需要mask,客户端发给服务端的帧必须mask,这既是协议规定,也是Netty编解码器里自动处理的事情。
推送场景里我们最常用的是文本帧,opcode为0x1。二进制帧0x2多用于文件、图片流;Ping/Pong帧0x9/0xA用于存活探测;关闭帧0x8用于主动断开。理解帧结构的意义在于:当你看到推送内容“变成乱码”或“只能收到第一帧”时,问题往往不在业务代码,而是编解码器没有按协议分帧。Netty的WebSocketServerProtocolHandler会帮你完成握手、帧解码和关闭帧响应,但你依然需要知道自己发出的数据最终是TextWebSocketFrame还是BinaryWebSocketFrame。
另一个容易忽略的点是WebSocket本身没有“主题”“路由”的概念。它只是一条全双工管道,服务端要把消息推给谁,完全靠业务代码维护“用户ID到Channel”的映射。SpringBoot内置的WebSocket走的是JSR-356标准,API里叫Session,Netty里叫Channel。无论哪种,协议层面都不会告诉你“这个用户关注了什么”,这些语义全靠自己在上层建模。理解这一点,后文讲用户与Channel的映射管理才有根基。
2.2 SpringBoot自带WebSocket为什么还要上Netty
SpringBoot的spring-boot-starter-websocket基于Tomcat或Jetty的WebSocket实现,API友好,配合@ServerEndpoint注解写起来非常快。但它的线程模型和Netty有本质区别。Tomcat的WebSocket连接由Tomcat自身的连接器管理,为每个连接分配的业务逻辑处理依赖容器线程池。连接少时没问题,一旦连接数上去,容器线程被慢业务占满,所有连接都会被拖慢,而且你很难精细控制每个连接的读写缓冲、心跳间隔和异常后的重连策略。
Netty则是一个独立的NIO网络框架,不依赖Servlet容器。它自己管理EventLoop线程,每个Channel绑定到一个EventLoop上,同一个Channel的所有读写事件都在同一个线程内串行执行,避免了锁竞争。SpringBoot在这个架构里退化为“业务容器”:用Spring管理对象、暴露HTTP接口,但真正的TCP/WebSocket接入由Netty监听端口完成。这样做的直接好处是连接生命周期、线程占用、背压处理都在你自己的代码里,出问题时能看穿每一层。
常见做法是让Netty绑定一个独立端口,比如8080给SpringBoot的HTTP接口,9090给WebSocket。这样既不用动原有接口,又能让长连接和短连接互不干扰。当然也可以让Netty复用SpringBoot的端口,但Config类里要处理端口占用冲突,而且HTTPS证书转发会多一层麻烦。我一般建议独立端口,尤其在生产环境,Nginx只需要把WebSocket的/ws路径代理到9090,HTTP请求照旧走8080,路由清晰,排查问题时少掉一半噪音。
2.3 Netty的ChannelPipeline与EventLoop模型:长连接推送的主心骨
Netty的每个Channel都绑定一条Pipeline,Pipeline里是一串ChannelHandler。入站消息从头部进,依次经过解码器、业务Handler;出站消息从尾部进,经过编码器后写到Socket。Push服务端的Pipeline我会这样排:
HttpServerCodec → HttpObjectAggregator → WebSocketServerProtocolHandler → TextWebSocketFrameHandler → 业务消息HandlerHttpServerCodec负责解析HTTP握手请求;HttpObjectAggregator把分片的HTTP消息聚合为完整请求;WebSocketServerProtocolHandler完成协议升级,之后自动移除HTTP编解码器,后续收到的都是WebSocketFrame;TextWebSocketFrameHandler里校验帧类型,把Text帧交给业务层处理。
EventLoop模型方面,Netty的Reactor主从结构里,bossGroup负责接受连接,workerGroup负责IO读写。每个worker线程维护一个Selector和一组Channel,读写事件在对应Channel的EventLoop上串行执行。这带来一个约束:不要在Handler里做耗时超过几十毫秒的阻塞操作,否则这个Channel读写的所有消息都会排队。常见做法是把业务操作提交给单独的线程池,或者用DefaultEventExecutorGroup给Handler指定独立的执行线程,避免卡住IO线程。这个点在后文避坑章节还会重点展开。
3. 搭一个能跑的Demo:SpringBoot+Netty实现WebSocket推送
3.1 工程准备与依赖:最小pom配置
用Maven建一个SpringBoot工程,Java版本选11或17都行,核心依赖只有两个:Netty和SpringBoot。SpringBoot的Web依赖只在需要提供HTTP API时用,如果推送入口完全是Netty,甚至可以不引入spring-boot-starter-web。不过为了演示“用SpringBoot管理Bean、用Netty接收连接”,我这里保留它。
<dependencies> <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> </dependencies>逻辑说明:netty-all不用细分模块,编译期省心;版本号建议固定一个你验证过的,不要用latest。spring-boot-starter-web会带入Tomcat,注意Tomcat的8080端口和Netty监听端口错开。
3.2 Netty服务端启动器:绑定端口与初始化Channel
写一个NettyServer组件,在SpringBoot启动完成后拉起Netty,监听9090端口。这里用@Component配合ApplicationRunner,或者直接在@PostConstruct里启动。生产环境推荐用ApplicationRunner,能保证Spring容器已经就绪,后面注入的Mapper、Service不会空指针。
@Component public class NettyWebSocketServer implements ApplicationRunner { private final EventLoopGroup bossGroup = new NioEventLoopGroup(1); private final EventLoopGroup workerGroup = new NioEventLoopGroup(); @Override public void run(ApplicationArguments args) { try { ServerBootstrap bootstrap = new ServerBootstrap(); bootstrap.group(bossGroup, workerGroup) .channel(NioServerSocketChannel.class) .childHandler(new ChannelInitializer<SocketChannel>() { @Override protected void initChannel(SocketChannel ch) { ChannelPipeline pipeline = ch.pipeline(); pipeline.addLast(new HttpServerCodec()); pipeline.addLast(new HttpObjectAggregator(65536)); pipeline.addLast(new WebSocketServerProtocolHandler("/ws")); pipeline.addLast(new TextWebSocketFrameHandler()); pipeline.addLast(new PushMessageHandler()); } }) .option(ChannelOption.SO_BACKLOG, 1024) .childOption(ChannelOption.TCP_NODELAY, true) .childOption(ChannelOption.SO_KEEPALIVE, true); ChannelFuture future = bootstrap.bind(9090).sync(); future.channel().closeFuture().sync(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } @PreDestroy public void stop() { bossGroup.shutdownGracefully(); workerGroup.shutdownGracefully(); } }参数说明:bossGroup线程数设为1即可,它只负责accept,NioEventLoopGroup(1)避免多余线程;workerGroup不设参数,默认按CPU核数乘2创建线程。SO_BACKLOG是操作系统层面的连接等待队列长度,设1024能应对突发的握手请求。TCP_NODELAY通知内核关闭Nagle算法,降低小消息的推送延迟,这对实时消息至关重要。SO_KEEPALIVE只是让TCP层保活探测,应用层心跳仍然要自己实现。
3.3 自定义WebSocket处理器:握手、文本帧与心跳
握手由WebSocketServerProtocolHandler完成,我们写的Handler主要处理连接建立、关闭、异常和业务帧。下面这个TextWebSocketFrameHandler继承SimpleChannelInboundHandler<TextWebSocketFrame>,只处理文本帧,其它类型帧由基类丢弃。
public class TextWebSocketFrameHandler extends SimpleChannelInboundHandler<TextWebSocketFrame> { private static final AttributeKey<String> USER_ID_KEY = AttributeKey.valueOf("userId"); @Override public void channelActive(ChannelHandlerContext ctx) { // 连接建立,但此时还没握手完成,不适合注册 } @Override public void channelInactive(ChannelHandlerContext ctx) { Channel channel = ctx.channel(); String userId = channel.attr(USER_ID_KEY).get(); if (userId != null) { SessionManager.remove(userId); } } @Override protected void channelRead0(ChannelHandlerContext ctx, TextWebSocketFrame frame) { String text = frame.text(); // 约定首次消息是{"type":"auth","userId":"123"} // 后续消息是{"type":"ping"}或者业务消息 JsonNode json = JsonUtil.parse(text); if ("auth".equals(json.get("type").asText())) { String userId = json.get("userId").asText(); ctx.channel().attr(USER_ID_KEY).set(userId); SessionManager.add(userId, ctx.channel()); ctx.channel().writeAndFlush(new TextWebSocketFrame("{\"type\":\"auth\",\"code\":200}")); } else if ("ping".equals(json.get("type").asText())) { ctx.channel().writeAndFlush(new TextWebSocketFrame("{\"type\":\"pong\"}")); } // 其它业务消息交给下一个Handler或丢弃 } @Override public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { cause.printStackTrace(); ctx.close(); } }逻辑说明:握手完成后Netty移除HTTP编解码器,后续收发都是WebSocketFrame。首条认证消息我们约定在业务层完成,这样服务端才知道这个Channel属于哪个用户。心跳用客户端主动发ping、服务端回pong的方式,比服务端轮询节省资源。AttributeKey是Netty给Channel挂数据的方式,等价于给会话塞了一个userId标签,断开清理时用得到。
3.4 从SpringBoot向指定会话推送:用户与Channel的映射管理
现在实现一个SessionManager,用ConcurrentHashMap保存userId到Channel的映射。这个类要被Spring容器管理,方便从HTTP接口或任务调度里调用推送方法。
@Component public class SessionManager { private final ConcurrentHashMap<String, Channel> channels = new ConcurrentHashMap<>(); public void add(String userId, Channel channel) { Channel old = channels.put(userId, channel); if (old != null && old.isActive()) { old.close(); } } public void remove(String userId) { channels.remove(userId); } public void push(String userId, String message) { Channel channel = channels.get(userId); if (channel != null && channel.isActive()) { channel.writeAndFlush(new TextWebSocketFrame(message)); } else { // 可以记录离线消息,或触发回调 } } public void broadcast(String message) { channels.forEach((userId, channel) -> { if (channel.isActive()) { channel.writeAndFlush(new TextWebSocketFrame(message)); } }); } }参数说明:ConcurrentHashMap保证多线程读写安全,因为Netty的worker线程和Spring业务线程都会访问这个Map。push前必须判断Channel是否active,因为TCP断开的通知可能有延迟。broadcast在Channel数量大时建议分批发送,避免一次循环写太多导致GC压力。实际项目中,推送接口是这样一个Controller:
@RestController public class PushController { private final SessionManager sessionManager; public PushController(SessionManager sessionManager) { this.sessionManager = sessionManager; } @PostMapping("/push/{userId}") public ResponseEntity<String> push(@PathVariable String userId, @RequestBody String message) { sessionManager.push(userId, message); return ResponseEntity.ok("pushed"); } }到这里,一个“客户端连Netty、服务端通过HTTP接口推消息”的最小闭环已经通了。你可以用浏览器控制台或者在线WebSocket工具连接ws://localhost:9090/ws,先发送认证消息,再用Postman调用推送接口,观察客户端能否收到。
4. 消息推送常见问题排查与避坑:从连不上到推不动的5个教训
4.1 握手403或者一直Pending:Tomcat的WebSocket与Netty冲突
现象:客户端连接ws://localhost:9090/ws时,SpringBoot的后台日志显示收到HTTP请求,但连接一直pending,最终报403;或者Netty控制台根本看不到任何日志。
原因:SpringBoot自带的WebSocket实现和Netty同时工作,端口被Tomcat占用。常见做法是在Netty配置里用了8080端口,但Tomcat也监听8080,结果握手请求被Tomcat吃掉,Tomcat找不到对应的@ServerEndpoint,返回403。
解决:确认Netty绑定的是独立端口,最稳的是在application.yml里显式声明:
server: port: 8080 netty: port: 9090然后在服务启动器里读取@Value("${netty.port}")。另外检查Nginx代理配置,如果走Nginx转发,需要设置proxy_set_header Upgrade $http_upgrade和proxy_set_header Connection "upgrade",否则WebSocket握手也会失败。这个坑排在第一,因为十个连不上有八个是端口或代理配置问题。
4.2 推送时NettyWorker线程被业务阻塞
现象:某个用户推送后,其他人也出现延迟;或者CPU不高但推送吞吐骤降。在Handler里打了日志,发现某个业务方法耗时几十秒,之后所有连接都卡了。
原因:我在第2章讲过,同一Channel的所有事件都在同一个EventLoop线程上执行。如果在channelRead0里直接调用了数据库查询、远程RPC或者Thread.sleep,这个EventLoop上的所有Channel都会被阻塞,等于一个慢操作拖垮了一个线程上的所有连接。
解决:把耗时操作提交到独立的业务线程池。一个简单做法是在Pipeline里给业务Handler指定独立的EventExecutor:
private static final DefaultEventExecutorGroup EVENT_EXECUTORS = new DefaultEventExecutorGroup(16); // initChannel中: pipeline.addLast(EVENT_EXECUTORS, new TextWebSocketFrameHandler());这样Handler里的代码跑在业务线程池,IO线程只负责编解码和转发。另一个选择是在channelRead0里用CompletableFuture.runAsync把耗时逻辑丢进去。注意writeAndFlush要回到Channel对应的EventLoop执行,否则会有线程安全风险,Netty允许在任意线程调用write,只是会内部切换线程,性能上略有损耗。
4.3 客户端断网后Channel还活着:心跳与失效清理
现象:服务端内存里SessionManager的Channel数量只增不减,客户端已经关闭了App或断网,但服务端还在向这些Channel推送,推送失败消息堆积,最终内存溢出。
原因:TCP断网时,服务端不会立刻收到FIN包。如果客户端没有发送关闭帧,而且网络设备还在转发或丢弃包,服务端需要等TCP超时才知道连接死了。应用层不主动探测,这个僵尸连接可能存活几分钟甚至几小时。
解决:启用Netty的IdleStateHandler心跳机制。在Pipeline中加一个:
pipeline.addLast(new IdleStateHandler(60, 0, 0)); // 60秒读空闲检测自定义一个IdleStateTriggerHandler,在userEventTriggered里接收IdleStateEvent.READER_IDLE,然后关闭该Channel,并从SessionManager移除。客户端侧也要有保活机制,定时发送ping。两手配合,服务端60秒内没有收到任何帧就判定失效,主动清理。心电间隔不能太短,否则移动网络频繁唤醒客户端会加剧耗电;也不宜太长,否则僵尸连接存活过久。一般建议60~120秒。
4.4 多实例部署后消息推给了错的机器
现象:服务从单机扩展到两台后,客户端连上机器A,但推送服务通过负载均衡调用到了一台不在这台机器上的实例,导致用户收不到消息。日志里SessionManager查询不到Channel。
原因:每台机器只维护自己的内存会话表。如果是无状态HTTP接口,负载均衡随便转发都没事;但WebSocket长连接是有状态的,连接建立在哪台机器,后续推送就必须绕到那台机器。
解决:常见做法是把Session信息放到Redis,Key是userId,Value是实例ID和Channel在本地表的引用。推送时先在Redis查实例ID,如果不在本机,就通过内部RPC或消息队列转发到对应实例。或者更简单一点,用带粘滞会话的负载均衡,客户端取不到新地址就先保持长连接不断开。但粘滞只能缓解,机器重启或缩容还需要重新分配。另一个思路是推消息时用Redis发布订阅广播,每台机器只推送本地持有的Channel,这样就不需要路由精准性。小型项目可以先用广播,连接数过万再考虑精准路由。
4.5 内存泄漏和Channel堆积:引用计数与关闭姿势
现象:服务跑几天后,堆内存使用率缓慢上升,Full GC频繁。Dump堆发现大量DefaultChannelPromise和ChannelOutboundBuffer实例。
原因:发送消息时调用了channel.write()而不是writeAndFlush(),消息留在缓冲区里没被写出;或者TextWebSocketFrame创建后没有释放。Netty的帧对象是引用计数的,如果只进站读数据但不释放,或者出站发送后不清理,会产生泄漏。还有一个常见操作是往已关闭的Channel里反复写入,异常被吞掉后,write返回的Future无人处理。
解决:养成三个习惯。第一,入站消息读完后,除非是简单的frame.text()转字符串后立即丢弃,否则记得ReferenceCountUtil.release(frame)。第二,调用write后必须调flush,或者直接用writeAndFlush。第三,所有writeAndFlush返回的ChannelFuture建议加监听器,日志记录失败:
channel.writeAndFlush(new TextWebSocketFrame(message)) .addListener((ChannelFutureListener) future -> { if (!future.isSuccess()) { log.warn("push failed for userId: {}, reason: {}", userId, future.cause().getMessage()); } });Netty自带的泄漏检测器在启动参数加-Dio.netty.leakDetection.level=paranoid,可以看到详细泄漏位置。但检测只能帮你发现,真正修复还要靠代码规范。
5. 让这个Demo撑住线上:会话维度推送的进阶设计
照着前面的代码,你已经能在开发环境跑通推送。但线上和Demo的区别,主要在于“推给谁”和“怎么推”的细节。这里分享一套我常用的会话维度推送落地方案:会话表升级为带版本号的路由表,推送入口统一走异步发送。
路由表除了userId、Channel,还要记录实例ID、连接建立时间、最近心跳时间。推送时先查路由表,如果Channel在本机就直发,不在本机就发一条内部消息到对应实例。消息体里带上userId和payload,目标实例收到后查自己的路由表完成推送。这套机制比Redis广播精准,也不会出现“每个实例都把消息发一遍”的浪费。
异步发送层用一个阻塞队列加工作线程池。HTTP接口收到推送请求后,只负责把消息放进队列,不等待发送结果。工作线程从队列取出消息,遍历目标用户列表,逐个调用writeAndFlush。这里要控制并发数,Netty的写入本身很快,但大量小消息会触发频繁的系统调用,批量发送是常用优化:把100条消息打包成一个List,循环发送,然后统一flush。实测单机长连接2万左右时,这种批量模式比单条sent的吞吐提升30%以上。
消息重试也很关键。客户端心跳超时的情况下,服务端判定离线后不要把消息直接丢弃,写入离线表,等客户端重连后拉取。在线时消息发出但ChannelFuture失败,同样入重试队列,最多重试3次,间隔指数退避。这些逻辑听着多,但每个模块都很短,做一个类里的两个队列就能完成。
最后提醒一句:别把Spring注入的Bean无脑变成Netty的静态工具。Netty线程访问SpringBean时要注意懒加载和代理问题,我习惯把SessionManager做成SpringBean,但Handler自己new,通过构造器传入Manager。这样Handler不是SpringBean也能用,而且不影响代理。
这套方案的维护性比直接用@ServerEndpoint高一个台阶。当某天产品说“要支持按部门推送、按标签推送”,你只需要改造路由表的查询维度,连接层完全不用动。希望帮到你。
本文还有配套的精品资源,点击获取