ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

SpringBoot整合Netty构建高并发客服系统:WebSocket通信实战与踩坑总结

SpringBoot整合Netty构建高并发客服系统:WebSocket通信实战与踩坑总结 做Web客服系统这几年我踩过的坑比写过的代码还多。今天就把这套基于SpringBoot Netty的在线客服聊天系统源码掰开揉碎讲清楚从架构选型到Netty通信的核心细节再到实际部署中那些文档里不会写的坑一次性聊透。这套方案适合正在做客服系统、IM系统、或者想搞明白Netty WebSocket通信的同学参考无论是自己学习还是直接拿去做二次开发都能少走不少弯路。1. 项目设计与技术选型为什么非要用Netty先聊聊技术选型。很多同学一看到“Web客服聊天”就觉得用SpringBoot自带的WebSocket就行了何必引入Netty这个大家伙。这个想法在并发量小的时候没什么问题但一旦进入生产环境面对几千上万的在线访客Tomcat自带的WebSocket实现就开始捉襟见肘了。我最早做客服系统用的就是SpringBoot内嵌Tomcat的WebSocket当时觉得反正Spring封装好了ServerEndpoint注解一标代码写起来挺爽。结果上线第二周就出问题高峰期连接数冲到8000多Tomcat的线程池直接被打满CPU飙到90%用户消息延迟从几十毫秒变成十几秒客服那边根本没法干活。后来痛定思痛把通信层整体迁移到Netty同样8000连接线程只用了十几个CPU稳定在20%以内这才算彻底解决问题。1.1 Netty相比传统WebSocket方案的优势在哪Netty之所以能扛住高并发连接核心在于它的IO模型和线程模型。传统Tomcat WebSocket是“一个连接一个线程”的BIO模型连接一多线程切换开销就成了噩梦。而Netty基于NIO用的是Reactor线程模型一个EventLoop线程可以同时处理成千上万个连接的读写事件连接数上来之后并不会线性消耗线程资源。再就是Netty对WebSocket协议的支持非常完善握手、帧编解码、ping/pong心跳、半包粘包处理这些Netty都内置了对应的Handler开箱即用。而且Netty的内存管理用的是池化DirectBuffer读写性能比堆内ByteBuffer要快不少这在高频消息推送场景下差距非常明显。1.2 系统整体架构设计这套客服系统的整体架构可以拆成几个核心部分业务后端SpringBoot提供REST API负责登录鉴权、会话查询、历史消息、访客信息管理等业务功能通信层Netty独立端口提供WebSocket服务负责维持与浏览器端的实时双向通信消息路由Netty收到消息后通过HTTP回调用SpringBoot进行业务处理或者直接投递给目标用户消息存储聊天记录异步写入数据库不阻塞主链路前端客户端访客端和客服工作台通过WebSocket连接通信服务这里有个关键设计思路Netty并不直接操作数据库它只负责维持连接和消息转发。收到消息后要么直接推送给目标连接要么通过回调接口交给SpringBoot做业务处理。这样做的好处是通信层和业务层解耦以后不管是把通信层单独拆成微服务还是换成其他技术栈业务层都不用动。2. SpringBoot工程搭建与核心依赖配置聊完架构直接上手写代码。这套系统的工程结构并不复杂核心是几个部分SpringBoot负责业务APINetty负责WebSocket通信两者通过端口区分互不影响。2.1 Maven依赖引入与实际配置先看pom.xml里的关键依赖parent groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-parent/artifactId version2.7.18/version relativePath/ /parent dependencies dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency dependency groupIdio.netty/groupId artifactIdnetty-all/artifactId version4.1.100.Final/version /dependency dependency groupIdcom.alibaba/groupId artifactIdfastjson/artifactId version2.0.32/version /dependency dependency groupIdorg.projectlombok/groupId artifactIdlombok/artifactId /dependency /dependenciesNetty版本这里特别注意一下Netty的版本迭代很快4.1.x是稳定分支。有些同学喜欢追新用4.2.x我个人建议生产环境还是用4.1.x社区生态和踩坑资料都更成熟遇到问题也更容易搜到解决方案。2.2 Netty服务端配置核心参数解析Netty服务端的初始化这段代码是整个通信层的基石Component public class NettyServer { private EventLoopGroup bossGroup; private EventLoopGroup workerGroup; private Channel serverChannel; PostConstruct public void start() throws InterruptedException { // bossGroup负责处理TCP连接建立线程数建议为1-2 bossGroup new NioEventLoopGroup(1); // workerGroup负责处理IO读写事件默认线程数是CPU核心数*2 workerGroup new NioEventLoopGroup(); ServerBootstrap bootstrap new ServerBootstrap(); bootstrap.group(bossGroup, workerGroup) .channel(NioServerSocketChannel.class) .option(ChannelOption.SO_BACKLOG, 1024) .option(ChannelOption.SO_REUSEADDR, true) .childOption(ChannelOption.TCP_NODELAY, true) .childOption(ChannelOption.SO_KEEPALIVE, true) .childHandler(new ChannelInitializerSocketChannel() { Override protected void initChannel(SocketChannel ch) { ChannelPipeline pipeline ch.pipeline(); // HTTP编解码器用于WebSocket握手 pipeline.addLast(new HttpServerCodec()); pipeline.addLast(new HttpObjectAggregator(65536)); // WebSocket协议处理器 pipeline.addLast(new WebSocketServerProtocolHandler(/ws)); // 自定义消息处理器 pipeline.addLast(new ChatMessageHandler()); } }); ChannelFuture future bootstrap.bind(9000).sync(); serverChannel future.channel(); } PreDestroy public void destroy() { if (serverChannel ! null) { serverChannel.close(); } if (bossGroup ! null) { bossGroup.shutdownGracefully(); } if (workerGroup ! null) { workerGroup.shutdownGracefully(); } } }这里有几个参数值得专门讲一下。SO_BACKLOG设置成1024表示操作系统底层accept队列的最大长度。如果系统短时间内涌入大量连接请求而Netty处理不过来这些请求会先排队。生产环境里这个值太小会导致连接被拒绝。TCP_NODELAY必须设置为true它关闭了Nagle算法。Nagle算法会把小数据包攒起来一起发减少网络包数量但对实时聊天来说就是灾难——每条消息都可能因为系统等着攒包而延迟几百毫秒。客服系统对延迟敏感必须关闭。SO_KEEPALIVE这个参数容易被误解。TCP层的keepalive默认两个小时才探测一次对即时通信来说太慢了所以真正保活还得靠应用层的心跳机制后面我会专门讲。2.3 WebSocket握手与协议处理细节WebSocketServerProtocolHandler(/ws)这个Handler做了三件事接收并校验客户端的HTTP升级请求完成从HTTP到WebSocket的握手以及处理WebSocket帧的编解码包括半包粘包。路径参数/ws表示客户端连接时走的路径比如 ws://ip:9000/ws。如果你的系统里需要区分不同业务比如客服端、访客端走不同的连接可以在握手阶段通过URL参数来区分然后把这些信息保存到Channel的attr属性里。这里有个细节HttpObjectAggregator(65536)是必须加的。因为WebSocket握手是一个HTTP请求这个请求可能被拆分成多个TCP包。HttpObjectAggregator的作用是把这些碎片拼成一个完整的FullHttpRequest否则握手会失败。这个值要大于最大的HTTP请求体积默认64KB足够用。3. Netty核心通信层实现要点搞定了服务端初始化接下来是通信层的关键实现。这一部分直接决定系统能不能稳定支撑高并发也是整个源码里最有技术含量的地方。3.1 在线用户管理ChannelGroup与ConcurrentHashMap结合Netty自身提供了ChannelGroup来管理连接但实际做客服系统时光有ChannelGroup不够因为我们需要对连接做更细粒度的分类管理。我的方案是用ConcurrentHashMap把用户标识映射到Channel再加一个ChannelGroup统一管理所有在线连接两者配合Component public class UserChannelManager { // 用户ID - Channel 的映射 private final ConcurrentHashMapString, Channel userChannelMap new ConcurrentHashMap(); // 状态标志 private final ConcurrentHashMapString, Integer userStatusMap new ConcurrentHashMap(); public void addUser(String userId, Channel channel) { userChannelMap.put(userId, channel); userStatusMap.put(userId, 1); } public void removeUser(String userId) { userChannelMap.remove(userId); userStatusMap.remove(userId); } public Channel getChannel(String userId) { return userChannelMap.get(userId); } public boolean isOnline(String userId) { return userChannelMap.containsKey(userId); } }为什么要同时用两个数据结构因为对ChannelGroup做全局广播非常方便比如系统公告要推给所有在线用户而ConcurrentHashMap则能够O(1)地找到目标用户的Channel实现定向推送。两者配合才能兼顾全局广播和点对点通信两种场景。关于用户标识的生成我的做法是访客在打开聊天窗口时先调用SpringBoot接口创建会话得到sessionId和userId然后带着这两个参数去连接WebSocket。在WebSocket握手阶段Netty会回调handlerAdded方法在这里从URL参数中取出userId并绑定到Channel上。3.2 消息处理器业务逻辑与协议处理的衔接消息处理是整个系统的核心枢纽。我自定义的ChatMessageHandler继承自SimpleChannelInboundHandler重写了几个关键方法Component ChannelHandler.Sharable public class ChatMessageHandler extends SimpleChannelInboundHandlerWebSocketFrame { Override public void channelActive(ChannelHandlerContext ctx) { System.out.println(客户端连接建立 ctx.channel().id()); } Override protected void channelRead0(ChannelHandlerContext ctx, WebSocketFrame frame) { if (frame instanceof TextWebSocketFrame) { String text ((TextWebSocketFrame) frame).text(); handleMessage(ctx, text); } else if (frame instanceof CloseWebSocketFrame) { ctx.close(); } } private void handleMessage(ChannelHandlerContext ctx, String text) { // 消息内容格式化成JSON处理 JSONObject msgObj JSONObject.parseObject(text); String type msgObj.getString(type); String fromUserId msgObj.getString(fromUserId); String toUserId msgObj.getString(toUserId); String content msgObj.getString(content); switch (type) { case CHAT: // 点对点聊天消息 sendToUser(toUserId, text); break; case ONLINE: // 用户上线通知 UserChannelManager.addUser(fromUserId, ctx.channel()); break; case HEARTBEAT: // 心跳响应 ctx.channel().writeAndFlush(new TextWebSocketFrame({\type\:\PONG\})); break; case READ: // 消息已读回执 break; default: break; } } Override public void channelInactive(ChannelHandlerContext ctx) { // 连接断开清理用户在线状态 String userId ctx.channel().attr(AttributeKey.valueOf(userId)).get(); if (userId ! null) { UserChannelManager.removeUser(userId); } ctx.channel().close(); } Override public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { cause.printStackTrace(); ctx.close(); } }这里有个踩过的坑必须说一下ChannelHandler.Sharable这个注解。如果不加多个ChannelPipeline里的Handler实例必须不同所以一般我们会new一个Handler加到pipeline里。但如果Handler是通过Spring注入的比如要用Autowired调用service层方法默认是单例的就必须加ChannelHandler.Sharable注解否则会报错。另外判断WebSocketFrame类型时一定要优先处理CloseWebSocketFrame。浏览器端关闭页面时Netty会收到这个帧如果你不主动关闭连接连接会一直挂着直到系统超时回收。线上出现大量僵尸连接的时候八成就是这个原因。3.3 心跳机制如何精准识别并清理死连接这是Netty通信里最容易出问题的地方。TCP连接断开时如果客户端异常崩溃比如断网、断电服务端并不能立刻感知到。如果没有心跳机制这些死连接会一直占用系统资源。Netty提供了IdleStateHandler来检测空闲连接用法很简单// 在ChannelInitializer中加一个空闲检测Handler pipeline.addLast(new IdleStateHandler(60, 0, 0, TimeUnit.SECONDS)); // 第一个参数读空闲第二个是写空闲第三个是总空闲 // 这里设置60秒内没有读到客户端数据就触发事件然后在ChatMessageHandler里重写userEventTriggered方法Override public void userEventTriggered(ChannelHandlerContext ctx, Object evt) throws Exception { if (evt instanceof IdleStateEvent) { // 超过60秒没有收到客户端数据判定为死连接 System.out.println(连接空闲超时关闭连接 ctx.channel().id()); ctx.close(); } else { super.userEventTriggered(ctx, evt); } }客户端这边要做的就是在Timeout之外主动发送心跳包。我通常要求前端每隔30秒发一次ping服务端60秒没收到就断开。这样两端配合既给了网络抖动预留一定时间又能在两分钟内清理掉死连接。这里再说一个细节把连接关闭的时候最好在channelInactive里再做一次用户下线处理。因为用户正常关闭页面时会走close流程但网络异常断开时不会。把清理逻辑写在channelInactive里无论是哪种情况都能兜底。4. 消息协议设计与业务流转逻辑通信层搭好之后剩下的重点就是业务层怎么跟通信层对接以及消息协议怎么设计。很多初级项目死就死在消息协议设计得一团糟后来维护成本居高不下。4.1 消息协议格式设计详解我用的消息协议是JSON格式通过type字段区分消息类型这样扩展性最好。下面是几个核心消息类型的定义ONLINE用户上线携带userIdOFFLINE用户下线CHAT聊天消息携带fromUserId、toUserId、content、timestampREAD消息已读回执SERVICE_REQUEST用户发起客服会话请求SERVICE_ASSIGN分配客服HEARTBEAT心跳包SYSTEM系统通知举个具体的聊天消息例子{ type: CHAT, fromUserId: USER_10001, toUserId: AGENT_20001, content: 你好我想咨询一下订单物流问题, timestamp: 1700000000000, msgId: MSG_20241101001 }msgId这个字段很容易被忽略但生产环境里它的作用非常大。因为网络不稳定时WebSocket消息可能会重发如果没有msgId做幂等处理用户就会看到两条一模一样的消息。我一般会用时间戳随机数或者UUID生成msgId服务端在存储消息时先检查msgId是否已存在重复消息直接丢弃。4.2 消息分发流程与线程模型Netty收到一条聊天消息后完整的处理流程是这样的Netty的IO线程从TCP缓冲区读取数据解码为WebSocketFrameChatMessageHandler解析消息内容判断类型如果是业务消息需要存储通过线程池异步提交给SpringBoot处理避免阻塞IO线程如果需要推送给其他在线用户直接通过UserChannelManager找到目标Channel并写入通知消息发送方消息已送达这里有一个特别重要的原则不要在Netty的EventLoop线程里执行耗时操作。Netty的IO线程是非常宝贵的资源它要处理成千上万个连接的读写。如果你在ChannelRead0里直接查数据库、调用远程接口一次操作几十毫秒这个线程上的其他所有连接都会被卡住。正确做法是把耗时操作丢给业务线程池Service public class MessageDispatchService { Autowired private ThreadPoolTaskExecutor taskExecutor; public void handleChatMessage(JSONObject msgObj) { // 异步保存消息记录 taskExecutor.execute(() - { // 调用SpringBoot的service层保存消息 messageService.saveMessage(msgObj); }); } }4.3 客服分配的几种策略与落地实现客服系统还有一个核心功能访客进来之后怎么分配一个客服给他。我试过几种策略各有优劣轮询分配按客服列表顺序轮流分配实现简单但不够智能有的客服能力强却分得少最少连接分配优先分配给当前会话数最少的客服负载相对均衡技能组匹配根据访客的问题类型分配到对应的客服技能组体验最好但实现复杂度高手动分配由管理员手动指定客服我的建议是第一版上线用最少连接分配算法简单效果也过得去。核心代码如下public Agent getLeastLoadedAgent() { return agentList.stream() .filter(agent - agent.isOnline()) .min(Comparator.comparingInt(agent - agent.getCurrentSessionCount())) .orElse(null); }等业务跑顺了再在消息里加入serviceType字段升级为技能组匹配方案从谁能接进化到谁更合适接。5. 前端接入要点与消息推送实践说完了后端聊聊前端怎么接入这套WebSocket通信服务。这个部分看似简单但实际开发中遇到的各种边界问题并不少。5.1 前端WebSocket连接与断线重连前端用浏览器原生WebSocket对象就够了不需要额外引库。但有两个关键点直接影响用户体验连接时机和断线重连。连接时机方面我建议页面加载后就创建连接但不要急着发消息。先调后端接口拿到userId和sessionId然后用拼接好的URL创建连接// 连接WebSocket服务 const wsUrl ws://10.0.0.5:9000/ws?userId${userId}sessionId${sessionId}; const socket new WebSocket(wsUrl); socket.onopen function() { console.log(WebSocket连接成功); // 发送上线通知 sendMessage({ type: ONLINE, userId: userId }); }; socket.onmessage function(event) { const msg JSON.parse(event.data); // 根据消息类型处理 handleMessage(msg); }; socket.onclose function() { console.log(连接断开); // 执行重连逻辑 reconnect(); };断线重连这里有个常见错误无限重连。如果服务端挂了前端就会陷入连接失败-重连-再失败-再重连的死循环既浪费资源体验也差。我的做法是设置最大重连次数比如5次每次重连的间隔递增1秒、2秒、4秒、8秒、16秒超过最大次数后提示用户刷新页面。另外要注意连接断开的瞬间可能会有消息丢失重连成功后要主动向服务端询问未读消息。5.2 消息推送时序与已读未读处理一个成熟的客服系统用户发完消息后最关心的两个反馈是消息发出去了没和客服看到了没。我在消息推送里分了两个阶段来反馈。用户在输入框点击发送时前端先把消息渲染到聊天窗口状态标记为发送中。服务端返回消息确认回执后前端再把状态改为已发送。客服这边收到消息并打开会话时前端发送READ消息服务端把该会话的消息都标记为已读再通知访客端把消息状态改为已读。这个设计在WebSocket里有几个小坑。一是顺序问题如果READ和CHAT消息在同一个连接上并发到达乱序处理会导致已读状态错乱。所以我在前端对消息做了序号管理收到消息时先放到本地队列里按序号处理后更新UI。6. 生产环境部署与常见问题排查代码写完部署上线才是真正的考验。这一节我重点讲生产环境里踩过的坑和排查经验。6.1 Netty线程数与内存参数调优先看线程模型配置。在Netty服务端bossGroup线程数根据CPU核数设置即可一般1到2个就够了。workerGroup的线程数选择8到16个是常见配置。但要注意线程数并非越多越好。Netty的每一个EventLoop都绑定一个Selector线程数过多线程切换开销反而会拖慢吞吐量。内存这块Netty用的是堆外内存DirectBuffer。在业务低峰期用jstat看内存使用情况时JVM堆内存看着不高但进程的内存占用却很高这就是DirectBuffer在起作用。如果出问题观察Netty的PoolArena内存使用情况会更有参考价值。另外通过jvm参数-XX:MaxDirectMemorySize可以限制堆外内存总量防止极端情况下内存失控。对于消息量大的场景还有一个参数值得调WebSocketServerProtocolHandler可以设置maxFramePayloadLength限制单条消息的体积。默认是65536字节64KB对文本聊天够用但如果以后要支持图片消息可能要调大到1MB。注意改这个参数时HttpObjectAggregator的大小也要同步调整。6.2 常见问题排查实录做这套系统的过程中我整理了高频遇到的问题如果你也在做类似的项目可以直接对照排查。第一个问题是连接建立成功后消息发不出去。这个问题我排查了很久最后发现是前端测试时没有等WebSocket状态变成OPEN就发消息了。浏览器控制台里onopen还没触发消息自然发不出去。处理方式是前端设置一个模块级变量isConnected在onopen回调里置为true发送消息前先判断状态如果未连接则缓存待发。第二个问题一段时间后连接自动断开。这个绝大多数是心跳超时导致的。如果你发现客户端还在正常使用但服务端误判为死连接先检查客户端的心跳是否按约定间隔发送。我在上面设置的服务端心跳超时是60秒如果前端因为页面切后台被浏览器冻结尤其是移动端定时器不再触发超时断开就必然发生了。处理办法是监听页面visibilitychange事件页面从后台切回前台时立即发送一次心跳同步最新状态。第三个问题多个服务实例部署时用户被路由到不同机器消息怎么互通。这就要把Netty服务做成无状态节点借助Redis发布订阅或者消息队列来广播消息。我在重构时引入了Redis的Pub/Sub机制Netty收到消息后先发布到Redis频道所有服务实例订阅频道并把消息推送到各自管理的连接上。这样即使某个实例挂了其他实例也能接管用户连接可用性大幅提升。第四个问题并发高峰期CPU飙高。CPU飙高通常有几个原因GC频繁、线程过多、业务逻辑阻塞了IO线程。先用jstack抓线程快照看是不是有大量线程阻塞在数据库连接等待上再用jstat看GC频率。如果发现老年代增长很快检查消息对象是否被某些全局集合持有引用导致无法被回收。我之前遇到过一次连接泄漏就是因为删除用户时没能从UserChannelManager里移除记录导致Map膨胀到几万个对象无法回收。6.3 部署架构与运维监控建议正式上线的部署架构我建议是SpringBoot业务服务与Netty通信服务分开部署通过内部端口通信。当然如果项目初期规模不大也可以合并为同一个进程用不同端口区分HTTP和WebSocket。运维监控方面Netty自身提供了很多可观测的指标。我建议把以下数据接入监控平台在线连接数区分访客和客服、每秒消息收发量、消息推送平均耗时、连接成功率、心跳超时断开数。前三个反映系统健康度后两个能帮助判断用户侧网络环境是否有问题。连接数监控我用了一个简单的计数器每次addUser加一、removeUser减一定期上报到Prometheus。最后分享一点实战心得这套系统从1.0到现在的版本前后迭代了大半年。回头看最让我觉得值得的做法不是在技术选型上有多激进而是从一开始就把通信层和业务层彻底解耦了。Netty管连接、SpringBoot管业务两边的演进互不拖累后面加功能、优化性能时都省了很多事。如果现在让我重做一遍客服系统我会优先在这几个方向提前规划一定要有消息幂等机制这事越晚做代价越大心跳和断线重连的机制一定在第一版就做好别指望上线后再补前端消息序号管理和本地队列设计看着不起眼但遇到消息乱序时你就知道它有多重要了。最后就是监控在线人数和消息量这两个指标一定要从第一天就开始记录不然出了性能问题你连排查的方向都没有。
返回列表