☰
基于 Netty 实现的 WebSocket 服务端
2026/10/9 1:45:09 网站建设 项目流程

目录

  • 一、短轮询(Short Polling)
  • 二、长轮询(Long Polling)
  • 三、WebSocket 长连接
  • 四、总结
  • 五、基于 Netty 实现的 WebSocket 服务端
    • 1、整体架构说明
    • 2、NettyWebSocketServer:WebSocket 服务端
      • 1. 线程组创建
      • 2. 启动方法
      • 3. 销毁方法
      • 4. 服务器启动逻辑
      • 5. 初始化管道(ChannelPipeline)
    • 3、NettyWebSocketServerHandler:消息与事件处理
      • 1. 连接事件处理
      • 2. 消息读取与业务处理
      • 3. 根据类型分发逻辑
    • 4、整体运行流程图
    • 5、总结
  • 六、附录代码

在许多业务场景中,服务端都需要主动向 Web 客户端推送消息,不仅仅限于 IM 通讯系统。
例如:

  • 小红点提醒
  • 新消息提示
  • 审批流通知

为了实现这些效果,通常会有几种服务端推送 Web 的方案。下面分别介绍三种常见方式:


一、短轮询(Short Polling)

原理

短轮询是指Web 端不断地以固定时间间隔向服务端发送 HTTP 请求。
如果服务端有新消息,就会在某次请求中返回。

示例:
阿斌之前做过的一个 OA 系统中,为了让用户实时收到审批流提醒或小红点提示,
客户端每秒向服务端发送一次请求,等待后端返回数据。

适用场景

  • 扫码登录:短时间内频繁查询二维码状态。
  • 小型 OA 系统:客户端数量不大,服务端压力较小。

缺点

  • 大量无效请求:大部分请求没有新消息返回,浪费服务器资源。
  • 服务端压力大:在高并发场景(如万人群聊)下,频繁请求会导致服务端难以承受。

实现思路

客户端(浏览器)以固定时间间隔(例如 1 秒或 5 秒)不断向服务端发起 HTTP 请求,检查是否有新消息。

前端示例(JavaScript)

// 每隔 5 秒向后端请求一次setInterval(()=>{fetch('/api/notice').then(res=>res.json()).then(data=>{if(data.hasMessage){console.log('📩 有新消息:',data.message);}else{console.log('无新消息');}}).catch(err=>console.error('请求错误:',err));},5000);

二、长轮询(Long Polling)

原理

长轮询是对短轮询的一种改进。
不同点在于:当请求没有新消息时,服务端不会立即返回,而是将请求“挂起(Hang)”一段时间。

  • 如果在这段时间内有新消息产生,服务端立即返回;
  • 如果超时仍无消息,再返回空响应;
  • Web 端再发起下一次请求。

因此,客户端的请求超时时间需要设置得更长一些。

优点

  • 相比短轮询,大幅减少了无效请求;
  • 降低了网络和服务器 QPS(每秒请求数);
  • 客户端功耗更低。

缺点

  • 仍有无效请求:若等待时间内无消息,仍需重新发起请求;
  • 服务端压力依然较大:虽然降低了入口请求频率,但悬挂的请求仍会占用线程或连接资源。
    例如:若有 1000 个请求在等待,就可能有 1000 个线程在轮询后端存储资源。

实现思路

客户端发起请求后,服务端不立即响应。
如果有新数据,立即返回;否则**挂起一段时间(例如 30 秒)**后再返回空结果。
客户端收到响应后再发起下一次请求。

前端示例(JavaScript)

functionlongPolling(){fetch('/api/long-poll').then(res=>res.json()).then(data=>{if(data.hasMessage){console.log('📢 新消息:',data.message);}else{console.log('⏳ 暂无消息');}// 继续下一轮请求longPolling();}).catch(err=>{console.error('连接错误:',err);// 等待 3 秒后重试setTimeout(longPolling,3000);});}// 启动长轮询longPolling();

三、WebSocket 长连接

原理

相比轮询方式,WebSocket 是一种真正的双向通信方案。
通过在客户端与服务端之间建立一个持久的 TCP/IP 长连接,
实现全双工(Full-Duplex)通信,即服务端可以主动向客户端推送数据。

优点

  • 实现真正的实时通信;
  • 省去了轮询带来的网络与性能损耗;
  • 更高效、更节能。

缺点

  • 实现复杂度较高:需要维护连接状态、心跳检测、断线重连等逻辑。

四、总结

推送方案优点缺点适用场景
短轮询实现简单,适合轻量场景请求频繁、性能浪费大扫码登录、小型系统
长轮询降低无效请求服务端压力仍大中等规模系统
WebSocket实时推送、性能高效实现复杂、需维护连接大型系统、IM、通知推送

五、基于 Netty 实现的 WebSocket 服务端

1、整体架构说明

项目主要分成两个类:

类名作用
NettyWebSocketServer启动一个 WebSocket 服务器(负责网络层)
NettyWebSocketServerHandler处理客户端消息和连接事件(负责业务层)

整个流程大致是:

启动服务器 → 客户端连接 → HTTP 升级为 WebSocket → 建立长连接 → 处理消息与心跳。


2、NettyWebSocketServer:WebSocket 服务端

1. 线程组创建

privateEventLoopGroupbossGroup=newNioEventLoopGroup(1);privateEventLoopGroupworkerGroup=newNioEventLoopGroup(NettyRuntime.availableProcessors());
  • bossGroup:只负责接收客户端连接请求。
  • workerGroup:负责处理连接的读写事件。
  • NioEventLoopGroup:基于 NIO 实现的事件循环线程池。

简单理解:

boss 是“门卫”,worker 是“工人”。


2. 启动方法

@PostConstructpublicvoidstart()throwsInterruptedException{run();}
  • @PostConstruct:表示当 Spring 容器启动时自动执行。
  • 调用run()启动 WebSocket 服务器。

3. 销毁方法

@PreDestroypublicvoiddestroy(){bossGroup.shutdownGracefully();workerGroup.shutdownGracefully();}
  • 在容器关闭时释放资源,优雅关闭线程组。

4. 服务器启动逻辑

ServerBootstrapserverBootstrap=newServerBootstrap();

ServerBootstrap是 Netty 启动服务器的引导类。

然后配置:

serverBootstrap.group(bossGroup,workerGroup).channel(NioServerSocketChannel.class).option(ChannelOption.SO_BACKLOG,128).option(ChannelOption.SO_KEEPALIVE,true).handler(newLoggingHandler(LogLevel.INFO))

含义:

  • channel(...):使用 NIO 的 TCP 通信通道;
  • SO_BACKLOG:允许的最大排队连接数;
  • SO_KEEPALIVE:启用 TCP 心跳保活机制;
  • LoggingHandler:打印连接日志,方便调试。

5. 初始化管道(ChannelPipeline)

.childHandler(newChannelInitializer<SocketChannel>(){@OverrideprotectedvoidinitChannel(SocketChannelch)throwsException{ChannelPipelinepipeline=ch.pipeline();...}});

这里是 WebSocket 的“核心逻辑链”,每一个连接都会被分配一个“管道(pipeline)”。

管道中添加的处理器如下:

处理器功能
IdleStateHandler(30,0,0)30秒内无读操作,则触发“读空闲”事件(可用来检测心跳超时)
HttpServerCodec()HTTP 编解码器(WebSocket 握手阶段需要 HTTP)
ChunkedWriteHandler()支持大数据流分块写入(比如文件传输)
HttpObjectAggregator(8192)将 HTTP 的分段消息聚合成完整请求(最大8KB)
WebSocketServerProtocolHandler("/")负责将 HTTP 升级为 WebSocket 协议,并保持长连接
NettyWebSocketServerHandler()自定义的业务逻辑处理类(见下)

最后一句:

serverBootstrap.bind(WEB_SOCKET_PORT).sync();

启动服务器,监听8090端口。


3、NettyWebSocketServerHandler:消息与事件处理

继承:

publicclassNettyWebSocketServerHandlerextendsSimpleChannelInboundHandler<TextWebSocketFrame>

TextWebSocketFrame表示 WebSocket 的文本帧(即发送的文本消息)。


1. 连接事件处理

@OverridepublicvoiduserEventTriggered(ChannelHandlerContextctx,Objectevt)

握手事件

if(evtinstanceofWebSocketServerProtocolHandler.HandshakeComplete){System.out.println("握手完成");}

当 HTTP 协议成功升级为 WebSocket 时打印“握手完成”。

心跳超时

elseif(evtinstanceofIdleStateEvent){if(event.state()==IdleState.READER_IDLE){System.out.println("读空闲");ctx.channel().close();}}
  • 超过 30 秒没收到消息(空闲),说明客户端断开或网络异常;
  • 服务器主动关闭连接。

2. 消息读取与业务处理

@OverrideprotectedvoidchannelRead0(ChannelHandlerContextctx,TextWebSocketFramemsg)

当客户端发来消息时执行。

Stringtext=msg.text();WSBaseReqwsBaseReq=JSONUtil.toBean(text,WSBaseReq.class);
  • 将客户端发送的 JSON 文本反序列化为WSBaseReq对象;
  • WSBaseReq里通常有一个字段表示消息类型type。

3. 根据类型分发逻辑

switch(WSReqTypeEnum.of(wsBaseReq.getType())){caseAUTHORIZE:break;caseHEARTBEAT:break;caseLOGIN:System.out.println("请求二维码");ctx.channel().writeAndFlush(newTextWebSocketFrame("123"));}
类型含义行为
AUTHORIZE授权(登录验证)暂未实现
HEARTBEAT心跳检测暂未实现
LOGIN登录请求(如扫码登录)打印“请求二维码”,并向客户端返回字符串“123”

4、整体运行流程图

浏览器(客户端) ↓ HTTP握手 NettyWebSocketServer ↓ 升级协议(状态码101) ↓ 建立长连接 NettyWebSocketServerHandler ↓ 监听消息(TextWebSocketFrame) ↓ 分发业务逻辑 ↓ 服务器可随时向客户端推送消息

5、总结

模块职责
NettyWebSocketServer搭建底层服务器、配置管道、维持长连接
NettyWebSocketServerHandler处理业务逻辑:握手、心跳、消息分发
IdleStateHandler心跳检测,断线清理
WebSocketServerProtocolHandlerHTTP → WebSocket 协议升级
TextWebSocketFrameWebSocket 的消息载体

六、附录代码

package com.donglin.mallchat.common.websocket;import io.netty.bootstrap.ServerBootstrap;import io.netty.channel.ChannelInitializer;import io.netty.channel.ChannelOption;import io.netty.channel.ChannelPipeline;import io.netty.channel.EventLoopGroup;import io.netty.channel.nio.NioEventLoopGroup;import io.netty.channel.socket.SocketChannel;import io.netty.channel.socket.nio.NioServerSocketChannel;import io.netty.handler.codec.http.HttpObjectAggregator;import io.netty.handler.codec.http.HttpServerCodec;import io.netty.handler.codec.http.websocketx.WebSocketServerProtocolHandler;import io.netty.handler.logging.LogLevel;import io.netty.handler.logging.LoggingHandler;import io.netty.handler.stream.ChunkedWriteHandler;import io.netty.handler.timeout.IdleStateHandler;import io.netty.util.NettyRuntime;import io.netty.util.concurrent.Future;import lombok.extern.slf4j.Slf4j;import org.springframework.context.annotation.Configuration;import javax.annotation.PostConstruct;import javax.annotation.PreDestroy;@Slf4j @ConfigurationpublicclassNettyWebSocketServer{publicstaticfinalintWEB_SOCKET_PORT=8090;// 创建线程池执行器privateEventLoopGroupbossGroup=newNioEventLoopGroup(1);privateEventLoopGroupworkerGroup=newNioEventLoopGroup(NettyRuntime.availableProcessors());/** * 启动 ws server * * @return * @throws InterruptedException */@PostConstructpublicvoidstart()throwsInterruptedException{run();}/** * 销毁 */@PreDestroypublicvoiddestroy(){Future<?>future=bossGroup.shutdownGracefully();Future<?>future1=workerGroup.shutdownGracefully();future.syncUninterruptibly();future1.syncUninterruptibly();log.info("关闭 ws server 成功");}publicvoidrun()throwsInterruptedException{// 服务器启动引导对象ServerBootstrapserverBootstrap=newServerBootstrap();serverBootstrap.group(bossGroup,workerGroup).channel(NioServerSocketChannel.class).option(ChannelOption.SO_BACKLOG,128).option(ChannelOption.SO_KEEPALIVE,true).handler(newLoggingHandler(LogLevel.INFO))// 为 bossGroup 添加 日志处理器.childHandler(newChannelInitializer<SocketChannel>(){@OverrideprotectedvoidinitChannel(SocketChannelsocketChannel)throwsException{ChannelPipelinepipeline=socketChannel.pipeline();//30秒客户端没有向服务器发送心跳则关闭连接pipeline.addLast(newIdleStateHandler(30,0,0));// 因为使用http协议,所以需要使用http的编码器,解码器pipeline.addLast(newHttpServerCodec());// 以块方式写,添加 chunkedWriter 处理器pipeline.addLast(newChunkedWriteHandler());/** * 说明: * 1. http数据在传输过程中是分段的,HttpObjectAggregator可以把多个段聚合起来; * 2. 这就是为什么当浏览器发送大量数据时,就会发出多次 http请求的原因 */pipeline.addLast(newHttpObjectAggregator(8192));//保存用户ip// pipeline.addLast(new HttpHeadersHandler());/** * 说明: * 1. 对于 WebSocket,它的数据是以帧frame 的形式传递的; * 2. 可以看到 WebSocketFrame 下面有6个子类 * 3. 浏览器发送请求时: ws://localhost:7000/hello 表示请求的uri * 4. WebSocketServerProtocolHandler 核心功能是把 http协议升级为 ws 协议,保持长连接; * 是通过一个状态码 101 来切换的 */pipeline.addLast(newWebSocketServerProtocolHandler("/"));// 自定义handler ,处理业务逻辑pipeline.addLast(newNettyWebSocketServerHandler());}});// 启动服务器,监听端口,阻塞直到启动成功serverBootstrap.bind(WEB_SOCKET_PORT).sync();}}
package com.donglin.mallchat.common.websocket;import cn.hutool.json.JSONUtil;import com.abin.mallchat.common.websocket.domain.enums.WSReqTypeEnum;import com.abin.mallchat.common.websocket.domain.vo.req.WSBaseReq;import io.netty.channel.ChannelHandlerContext;import io.netty.channel.SimpleChannelInboundHandler;import io.netty.handler.codec.http.websocketx.TextWebSocketFrame;import io.netty.handler.codec.http.websocketx.WebSocketServerProtocolHandler;import io.netty.handler.timeout.IdleState;import io.netty.handler.timeout.IdleStateEvent;/** * Description: * Author: <a href="https://github.com/zongzibinbin">abin</a> * Date: 2023-08-27 */publicclassNettyWebSocketServerHandlerextendsSimpleChannelInboundHandler<TextWebSocketFrame>{@OverridepublicvoiduserEventTriggered(ChannelHandlerContextctx,Objectevt)throwsException{if(evt instanceof WebSocketServerProtocolHandler.HandshakeComplete){System.out.println("握手完成");}elseif(evtinstanceofIdleStateEvent){IdleStateEventevent=(IdleStateEvent)evt;if(event.state()==IdleState.READER_IDLE){System.out.println("读空闲");//todo 用户下线ctx.channel().close();}}}@OverrideprotectedvoidchannelRead0(ChannelHandlerContextctx,TextWebSocketFramemsg)throwsException{Stringtext=msg.text();WSBaseReqwsBaseReq=JSONUtil.toBean(text,WSBaseReq.class);switch(WSReqTypeEnum.of(wsBaseReq.getType())){caseAUTHORIZE:break;caseHEARTBEAT:break;caseLOGIN:System.out.println("请求二维码");ctx.channel().writeAndFlush(newTextWebSocketFrame("123"));}}}

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询