ARTICLE DETAIL

资讯详情

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

手把手搭建直播高并发后端:从Netty长连接到消息队列实战

手把手搭建直播高并发后端:从Netty长连接到消息队列实战 “后端小白”这个标签我在自己博客的自我介绍里挂了快一年直到上个月我终于用一个周末把一套直播场景的高并发后端环境从零跑通了。说实话这套环境放到大厂里也就是个玩具但对于一个平时只写CRUD、连消息队列都没在生产环境摸过的人来说能把直播间在线人数、弹幕消息、礼物通知这些高频交互全部扛下来那一刻是真的爽。这篇文章就把我从选型、编码、压测到优化的完整过程整理出来全程是“小白能直接抄作业”的颗粒度。如果你也好奇直播这类高并发IM场景的后端到底长什么样或者你正站在从“增删改查”往后端架构进阶的路口这篇笔记应该能帮你省掉不少弯路。1. 为什么一个后端小白要“手搓”直播高并发环境1.1 直播后端到底在解决什么问题我最初接触到直播需求的时候第一反应是“直播后端不就是给App提供几个接口吗拉流地址、用户信息、礼物列表完事了”。真正开始设计之后才发现直播场景和传统后台系统完全是两个物种。传统CRUD系统里用户主动发一次请求服务器处理完返回一次响应连接就结束了。哪怕做成 RESTful API一次会话也就几秒钟一秒钟能支撑几百个请求已经算不错了。但直播场景完全不同观众进入直播间之后需要和服务器维持一条长连接服务器要持续不断地把弹幕、礼物、进场通知、系统公告推送给所有在线用户。这意味着连接不是“用完就断”的而是要长期驻留的。服务器不再是“被动响应”而是需要“主动推送”。高并发不再是“一个瞬间来了很多请求”而是“一个直播间里同时挂了几万条活跃连接”。用一个生活化的类比来理解传统后端像银行的柜台窗口用户排队来办业务办完就走了柜员服务完一个人再接下一个直播后端更像一个大礼堂里的中央空调所有人在里面待着不走空调得持续给整个空间供应冷气人越多空调的负载越大而且任何一个角落缺了冷气都会有人喊热。所以做直播后端核心不是“接口性能”而是连接管理和消息分发。为什么我之前看的很多后端招聘JD里会强调“IM系统”“长连接”“消息推送”经验因为直播、聊天室、联机游戏、协同办公软件底层都是同一套东西。这也是这次项目里我选 WebSocket 和 Netty 作为核心技术栈的直接原因。1.2 目标定义什么才算“搭起来”目标如果定得太模糊很容易做着做着就跑偏。我在动手之前给自己定了三个非常具体的验收标准你也可以按这个思路来定义自己的项目第一单台服务器能支撑10000 个并发在线连接。这是对标一个小型直播间或中型聊天室的实际数据如果单台连一万连接都扛不住后面谈优化就没有意义。第二在弹幕高峰期一条消息从发出到其他用户收到延迟要控制在1 秒以内。弹幕的实时性要求比普通聊天更高如果延迟超过一秒用户会明显感觉到“卡”。第三有一个简单的运维看板能实时看到当前在线人数、消息吞吐量和服务器负载。这三个指标立下来之后整个项目的边界就清楚了我需要一套支持长连接、能广播消息、可观测运行状态的后端系统。这也决定了后面的技术选型——不是选“最流行”或“最热门”的框架而是选“在最坏情况下也能扛住上述指标”的组件组合。2. 搭建前必须先想清楚的三件事2.1 技术栈选型为什么选 Spring Boot Netty 而不是只用 Spring MVC我在网上查资料的时候发现一个很有意思的现象很多直播后端的技术分享会提到 Netty但 Spring Boot 自带的 WebSocket 好像也够用那到底选哪个先说结论如果你想快速做个 DemoSpring Boot 内置 WebSocket 完全够用但如果你要的是高并发长连接场景下的稳定性和可控性Netty 是更合理的选择。为什么要二选一而不是只用一个因为这两个东西的定位不同。Spring Boot 是应用开发框架负责处理业务逻辑、对接数据库、提供 REST APINetty 是高性能网络通信框架专门解决海量长连接的并发读写问题。直播场景里高并发的瓶颈在“连接层”和“消息收发层”这恰好是 Netty 的主场。我用Netty的时候它内置的EventLoop模型能在一个线程里处理成千上万个连接的 IO 事件线程资源占用极小而 Spring Boot 默认的 Tomcat 容器本质上还是为“短请求、短连接”设计的虽然也能跑 WebSocket但在连接数上去之后内存和线程的消耗会明显增大。为了验证这个判断我专门做了个小对比实验用同一台 8C16G 的测试服务器尝试建立 1 万个 WebSocket 连接。Tomcat 原生 WebSocket 能扛住但资源消耗接近 80%响应已经开始变慢Netty 在相同连接数下 CPU 和内存表现都从容很多。这个实验让我彻底下定了决心连接层用 Netty业务层用 Spring Boot两者各管各的中间通过消息队列解耦。至于这俩怎么整合第三节会详细说。2.2 高并发消息模型为什么说直播弹幕本质是 IM做系统设计的第一步不是写代码而是搞清楚数据流。直播场景里最常见的消息包括弹幕消息、礼物通知、进入直播间、点赞、系统公告。这些消息表面上五花八门但抽象来看都指向同一种架构模型IMInstant Messaging即时通讯模型。IM 模型的核心特征有两个。第一消息是实时双向的用户发弹幕服务器要把消息广播给房间里其他所有用户用户收到礼物通知同样需要服务器推送。第二消息有频道Channel的概念一条弹幕只应该发给当前直播间的观众不能发到其他直播间去。这就像微信群和QQ群的区别——你在 A 群发消息B 群的人不该收到。想通了这一点整个系统就简化为三个模块接入层负责维护客户端与服务器之间的长连接处理连接建立、鉴权、心跳和断线重连。路由层负责根据消息的直播间 ID把消息分发到对应的频道广播给该频道的所有连接。存储层负责把需要持久化的消息比如弹幕记录、礼物记录异步写入数据库供后续查询或回放。这三个模块各做各的事彼此不纠缠。当你把直播后端理解成“多房间 IM 系统”之后再看网上那些高并发 IM 的教程和开源项目会发现很多东西都是相通的学习成本一下子降下来不少。2.3 最小闭环单机架构能不能撑起直播场景很多人一提高并发就想到微服务、Kafka、Redis Cluster、K8s。但我要说的是从一个后端的角度出发如果单机都跑不通分布式就是空中楼阁。所以这次搭建我给自己定下的原则是——先跑通单机最小闭环再谈扩展。最小闭环包括哪些东西进程层面一个应用Spring Boot 工程内部嵌入了 Netty 来托管 WebSocket 连接同时提供了 REST API 供客户端查询直播间状态。中间件层面一台 Redis 负责在线人数统计和热点缓存一台 RabbitMQ 负责弹幕消息的异步削峰和落库。数据库层面一台 MySQL 用来存弹幕记录和礼物流水。这些组件全部可以通过 Docker Compose 在一台服务器上跑起来。有人会问只用一台服务器这算不算“玩具”我的看法是单机架构跑出来的数据和踩出来的坑比看十篇分布式架构文章都管用。比如当你压测到 1 万连接时你会真真切切地感受到 Linux 文件描述符限制是什么Netty 的堆外内存是怎么增长的RabbitMQ 队列堆积会对延迟造成什么影响。没有这些体感直接去搭 K8s 集群你连问题出在哪一层都不知道。单机闭环搭通之后再往多节点演进反而是顺理成章的事前面加 Nginx 做负载均衡把连接分散到多台实例上连接层和业务层之间用 Redis Pub/Sub 或者消息队列做广播转发。到那时候你已经有“体感”了知道瓶颈在哪知道该监控什么参数扩容的每一步都有数据支撑。3. 从 0 到 1 的完整搭建过程3.1 初始化工程与基础配置我这次选的框架版本是 Spring Boot 3.2.x JDK 17。选 Spring Boot 3 而不是 2.x主要是看中它内置支持 Netty 相关的依赖管理更干净同时性能比 2.x 有提升。当然如果你还在用 JDK 8用 Spring Boot 2.7 也完全能走通下面的流程Netty 的 API 差异不大。工程结构长这样live-highconcurrency-backend ├── src/main/java/com/example/live │ ├── config # 配置类 │ ├── netty # Netty 服务端、Handler、协议编码器 │ ├── controller # REST API 控制器 │ ├── service # 业务逻辑层 │ ├── dao # 数据库访问层 │ ├── entity # 实体类 │ ├── mq # RabbitMQ 生产者/消费者 │ └── LiveApplication.java ├── src/main/resources │ ├── application.yml │ └── logback-spring.xml ├── docker-compose.yml # Redis RabbitMQ MySQL └── pom.xmlPOM 里的关键依赖如下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 groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-data-redis/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-amqp/artifactId /dependency dependency groupIdcom.baomidou/groupId artifactIdmybatis-plus-boot-starter/artifactId version3.5.5/version /dependency dependency groupIdmysql/groupId artifactIdmysql-connector-j/artifactId scoperuntime/scope /dependency这里有个容易忽略的点Spring Boot 3.x 下 MyBatis-Plus 需要选择支持jakarta命名空间的版本3.5.3否则启动就会报错。我第一次搭建时就是在这里卡了半个小时。application.yml里除了常规的数据源和 Redis 配置还要单独给 Netty 设置端口和参数server: port: 8080 netty: port: 9090 boss-threads: 2 worker-threads: 8 so-backlog: 1024 write-buffer-high-water-mark: 65536 write-buffer-low-water-mark: 32768boss-threads负责接受新连接一般设 1 到 2 个就行worker-threads负责处理 IO 事件理想情况下等于 CPU 核数或者稍多一点。我机器是 8 核所以设了 8 个。write-buffer-high-water-mark和write-buffer-low-water-mark是 Netty 写出缓冲区的水位线这个参数直接影响高并发下消息广播的内存占用后面踩坑部分会细说。中间件直接用一个docker-compose.yml拉起来version: 3 services: mysql: image: mysql:8.0 environment: MYSQL_ROOT_PASSWORD: root123 MYSQL_DATABASE: live_db ports: - 3306:3306 volumes: - mysql-data:/var/lib/mysql redis: image: redis:7.0 ports: - 6379:6379 rabbitmq: image: rabbitmq:3.12-management ports: - 5672:5672 - 15672:15672 volumes: mysql-data:我在实际开发中更习惯把 Netty 服务端通过PostConstruct在应用启动后自动拉起这样 Spring Boot 的 IOC 容器已经初始化完毕Netty 的 Handler 里可以直接注入 Service 和 RedisTemplate省去一堆手动维护 Bean 的麻烦。3.2 直播房间连接管理从用户进场到挥手告别Netty 服务端的核心启动类写完以后接下来是连接接入的逻辑。这一块是整个系统里最基础也最容易出错的地方。先看服务端初始化代码Component public class WebSocketServer { private final ChannelGroup roomGroup new DefaultChannelGroup(GlobalEventExecutor.INSTANCE); // 房间维度roomId - ChannelGroup private final ConcurrentHashMapString, ChannelGroup roomChannels new ConcurrentHashMap(); PostConstruct public void start() throws InterruptedException { EventLoopGroup bossGroup new NioEventLoopGroup(nettyProps.getBossThreads()); EventLoopGroup workerGroup new NioEventLoopGroup(nettyProps.getWorkerThreads()); try { ServerBootstrap bootstrap new ServerBootstrap(); bootstrap.group(bossGroup, workerGroup) .channel(NioServerSocketChannel.class) .option(ChannelOption.SO_BACKLOG, nettyProps.getSoBacklog()) .childOption(ChannelOption.SO_KEEPALIVE, true) .childOption(ChannelOption.TCP_NODELAY, true) .childOption(ChannelOption.WRITE_BUFFER_WATER_MARK, new WriteBufferWaterMark( nettyProps.getWriteBufferLowWaterMark(), nettyProps.getWriteBufferHighWaterMark())) .childHandler(new WebSocketServerInitializer(roomChannels)); ChannelFuture future bootstrap.bind(nettyProps.getPort()).sync(); log.info(Netty WebSocket Server started on port {}, nettyProps.getPort()); future.channel().closeFuture().sync(); } finally { bossGroup.shutdownGracefully(); workerGroup.shutdownGracefully(); } } }这里我用roomChannels这个ConcurrentHashMap来按直播间维度维护连接组。为什么要这样设计因为直播消息的广播范围是“房间级”的——进入房间 1001 的观众只收 1001 房间的弹幕。如果只用一个全局ChannelGroup每条消息都得遍历所有连接然后判断要不要发连接数上来后这个判断成本非常可观。接下来是关键一条新连接进来之后先做协议升级和参数解析。用户在建立 WebSocket 连接时需要在 URL 上带上roomId和userId比如ws://101.35.200.1:9090/ws?roomId1001userId9527。这样服务端在握手阶段就能知道这条连接属于哪个直播间直接把它加入对应的ChannelGroup。public class WebSocketServerInitializer extends ChannelInitializerSocketChannel { Override protected void initChannel(SocketChannel ch) { ChannelPipeline pipeline ch.pipeline(); pipeline.addLast(new HttpServerCodec()); pipeline.addLast(new HttpObjectAggregator(65536)); pipeline.addLast(new WebSocketServerCompressionHandler()); pipeline.addLast(new WebSocketServerProtocolHandler(/ws)); pipeline.addLast(new WebSocketFrameHandler(roomChannels)); } }在WebSocketFrameHandler的channelRead0方法里处理几种常见的请求用户进入房间从 URI 解析 roomId 和 userId加入roomChannels对应的 ChannelGroup把当前在线人数广播给该房间所有人。心跳消息用户端每隔 30 秒发送一条{type: ping}服务端收到后原路返回{type: pong}用于维持连接不被中间网络设备断开。弹幕消息收到{type: chat, roomId: 1001, userId: 9527, content: 大家好}之后先写入 Redis 缓存再推送给当前 roomId 对应 ChannelGroup 里的所有连接。断开连接在channelInactive中把 Channel 从所属房间的 ChannelGroup 里移除并广播离场通知防止后续广播消息发给一个已经失效的 Channel。这里有个容易踩的坑不要直接在channelRead0里同步调用 Redis 或数据库。为什么因为 Netty 的 IO 线程是非常宝贵的资源所有连接的消息收发都依赖这几个线程一旦某个操作阻塞影响的是整个服务端的所有连接。正确做法是把消息丢给异步线程池或者通过消息队列处理后端逻辑。我当时在开发阶段没注意直接同步 Redis结果压测到 3000 连接时 CPU 被打满了排查下来发现是 Redis 连接等待把 IO 线程堵死了。3.3 消息广播一条弹幕是怎么送到所有人手里的很多人第一次接触 Netty 广播时会觉得“不就是循环遍历发送吗有什么难的”。事实上这条链路里藏着好几个性能杀手。先看朴素版本是怎么写的public void broadcast(String roomId, String message) { ChannelGroup group roomChannels.get(roomId); if (group ! null) { group.writeAndFlush(new TextWebSocketFrame(message)); } }writeAndFlush方法会遍历 ChannelGroup 里的所有 Channel然后逐个写入并发送。这个逻辑本身没问题但到了高并发场景就会出现两个问题。第一个问题是部分 Channel 写入速度不一致导致的内存堆积。每个客户端的网络情况不同有的客户端网速快消息秒收有的客户端网络差写入缓冲区的数据一直发不出去。如果服务端不管不问数据会在内存里越积越多最终导致内存溢出。这就是我配置write-buffer-high-water-mark和write-buffer-low-water-mark的原因——当某个 Channel 的待发送数据超过高水位线时Netty 会将该 Channel 标记为不可写状态低于低水位线后再恢复可写。我写的广播逻辑里加了一个判断public void broadcast(String roomId, String message) { ChannelGroup group roomChannels.get(roomId); if (group null || group.isEmpty()) { return; } String safeMessage message; for (Channel channel : group) { if (channel.isActive() channel.isWritable()) { channel.writeAndFlush(new TextWebSocketFrame(safeMessage)); } } }第二个问题是消息内容序列化开销。直播场景的弹幕消息非常频繁如果每次广播都把对象序列化成 JSON 字符串GC垃圾回收压力会非常大。后来我把消息对象提前序列化好只保留一个字节数组广播时直接构造TextWebSocketFrame减少重复的序列化开销。这算是一个小幅优化但在高并发下效果肉眼可见。广播模块完整流程可以参考这张表步骤操作说明1客户端发送消息通过 WebSocket 连接发送 JSON 消息2Netty 处理器解析校验消息格式提取 type、roomId、content 等字段3写入 Redis弹幕内容写入 Redis List用于读取最近消息、热度统计4投递消息队列生产者将消息发到 RabbitMQ消费端异步落 MySQL5房间内广播遍历房间对应的 ChannelGroup写入所有活跃连接6返回 ACK 给发送者发送者本人可以不通过广播收到消息而是立即回显减少一次消息延迟这里有个小设计为什么发送者本人不参与广播而是返回 ACK 回显因为一条弹幕走完整个广播流程后发送者的消息其实已经到了其他客户端那边如果发送者再从广播里收一遍不仅浪费带宽而且可能在弱网环境下造成“自己发的消息比别人慢”的错觉。所以我会对发送者单独回复一条 ACK 确认同时广播给房间内其他观众。这种设计在游戏和直播场景里很常见。3.4 压测验证8C16G 机器到底能扛多少连接代码写完、服务启动成功之后真正的考验才刚刚开始——压测。我一开始试图用 Postman 或浏览器控制台的 WebSocket 客户端来做压测但很快发现这完全是杯水车薪因为浏览器对并发连接数有严格的限制通常一个域名下最多 6 个 HTTP 连接WebSocket 虽然不受这个限制但浏览器本身也没法稳定压出几千个并发连接。所以我直接用 Java 写了一个模拟客户端用 Netty 分别创建 1 万个连接连到服务端每个连接建立成功后定时发送心跳和弹幕消息。压测脚本的关键代码大致如下public class WebSocketLoadClient { public static void main(String[] args) throws Exception { int totalConnections 10000; String serverUrl ws://101.35.200.1:9090/ws?roomId1001userId; EventLoopGroup group new NioEventLoopGroup(8); CountDownLatch latch new CountDownLatch(totalConnections); AtomicInteger connected new AtomicInteger(0); for (int i 0; i totalConnections; i) { String userId u_ i; WebSocketClient client new WebSocketClient( URI.create(serverUrl userId), group) { Override protected void onOpen() { int cnt connected.incrementAndGet(); if (cnt % 1000 0) { System.out.println(已连接: cnt); } // 连接建立后随机发送一条消息 send({\type\:\chat\,\roomId\:\1001\,\userId\:\ userId \,\content\:\hello\}); latch.countDown(); } }; client.connect(); } latch.await(); System.out.println(全部连接建立完成继续维持连接 5 分钟...); Thread.sleep(5 * 60 * 1000L); group.shutdownGracefully(); } }压测结果让我挺意外的在 8C16G 的云服务器上Netty 服务端平稳接收了 10000 个连接CPU 占用率在 70% 左右内存消耗接近 3.2G消息广播延迟基本在 100ms 以内。内存比预期高一些仔细排查后发现是每连接设置的心跳定时任务占用了不少内存每个定时任务约 30KB。我对心跳任务做了一层优化改成只用一个全局的调度任务轮询所有连接内存直接降了接近 1G。这个优化点值得记住能用公共调度解决的事情别为每个连接创建独立的定时任务。压测过程中我也记录了服务端的核心指标Netty 事件循环线程的pendingTasks、Channel 的写入水位、GC 频率。这些指标后来成了我调优和排查问题的重要依据。4. 高并发优化实践从“能跑”到“扛得住”4.1 接入层优化Nginx 反向代理 WebSocket单机跑通只是第一步要让它看起来像一个正经的高并发架构接入层必须有一个统一的流量入口。我选了 Nginx原因有三个一是配置简单二是性能好三是几乎零成本就能实现负载均衡和端口转发。但 Nginx 默认是不支持 WebSocket 的因为 WebSocket 协议握手时要进行 HTTP Upgrade升级之后连接从 HTTP 变成全双工通信后续不再走 HTTP 协议。这意味着 Nginx 需要特殊配置否则会直接断开 WebSocket 连接。核心配置段如下map $http_upgrade $connection_upgrade { default upgrade; close; } upstream live_backend { server 127.0.0.1:9090; keepalive 64; } server { listen 80; server_name live.example.com; location /ws { proxy_pass http://live_backend; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection $connection_upgrade; proxy_set_header Host $host; proxy_read_timeout 3600s; proxy_send_timeout 3600s; } location / { proxy_pass http://127.0.0.1:8080; proxy_set_header Host $host; proxy_set_header X-Real-IP $remote_addr; } }几个关键点需要解释map块的作用是把Upgrade请求头映射为对应的Connection头从而让 Nginx 正确识别 WebSocket 升级请求。proxy_read_timeout和proxy_send_timeout是 Nginx 在 WebSocket 场景下最重要的参数因为 WebSocket 连接是长连接如果这两个超时时间设置过短默认 60 秒连接空闲超过 60 秒就会被 Nginx 切断客户端全部掉线。我设成了 3600 秒再配合客户端 30 秒一次的心跳基本能保证连接稳定。keepalive 64的作用是让 Nginx 与后端 Netty 服务之间保持长连接复用。如果没有这个配置Nginx 每次转发请求都要新建一条 TCP 连接到后端高并发下会产生大量 TIME_WAIT 连接把服务器端口资源耗尽。加上keepalive后Nginx 直接复用已有连接性能提升很明显。我实测了一下不加 Nginx 时直连 Netty 和加了 Nginx 之后通过代理访问在 1 万连接场景下消息延迟基本没有影响但 Nginx 进程本身的 CPU 占用可能还需要单独观察。在后续多节点部署时只需要在 upstream 里加服务器地址就能实现连接级负载均衡这是目前最平滑的扩展方式。4.2 业务层优化消息队列削峰与异步落库直播弹幕有一个非常明显的流量特征瞬时洪峰。一场热门直播里主播说了一句劲爆的话可能在几秒内涌进来几千条弹幕。如果每条弹幕都直接写数据库MySQL 会当场被打爆。解决思路是引入消息队列做“削峰填谷”。我把消息队列的引入时机放在压测之后。为什么要压测之后再引因为如果你连最基础的广播链路都没跑通直接上 MQ 只会让问题更多排查起来没有任何基线可以对照。在 RabbitMQ 里的核心配置如下Configuration public class RabbitConfig { public static final String EXCHANGE live.exchange; public static final String QUEUE live.chat.queue; public static final String ROUTING_KEY live.chat; Bean public Queue chatQueue() { return QueueBuilder.durable(QUEUE).build(); } Bean public DirectExchange liveExchange() { return new DirectExchange(EXCHANGE); } Bean public Binding binding(Queue chatQueue, DirectExchange liveExchange) { return BindingBuilder.bind(chatQueue).to(liveExchange).with(ROUTING_KEY); } }生产者投递消息的代码很简单核心是设置消息的过期时间TTL和确认模式public void sendToQueue(RedPacketMessage message) { rabbitTemplate.convertAndSend(RabbitConfig.EXCHANGE, RabbitConfig.ROUTING_KEY, message, msg - { msg.getMessageProperties().setExpiration(5000); return msg; }); }消费端负责把消息批量插入 MySQL。这里同样有一个优化技巧不要来一条插一条而是攒一批再批量写入。I/O 开销是数据库写入的最大瓶颈批量插入可以将性能提升好几倍。我的消费端代码里维护了一个LinkedBlockingQueue等队列攒够 500 条或超过 1 秒就统一INSERT一次。这里要提醒一个小白容易犯的错不要把消息队列既当削峰工具又当数据存储。RabbitMQ 里的消息如果没有被消费端确认会一直保存在内存或磁盘中消息量大的时候很容易把中间件本身的内存吃光。所以消费端一定要开启手动 ACK且在消息处理完成后才能确认确保消息不会丢失也不会积压。我当时在开发阶段偷懒用了自动 ACK结果压测时队列积压了几万条消息RabbitMQ 的内存报警直接黄了。4.3 数据层优化Redis 缓存在线人数与热点数据直播场景里除了弹幕还有一个高频请求就是“在线人数”和“热度值”。这两个数据如果每次都从 MySQL 里实时聚合统计系统的 QPS 很快就会被拖垮。我在这个环节用了两层缓存策略。第一层是用 Redis 的HyperLogLog来做 UV 统计。HyperLogLog 是一种概率数据结构能用极少的内存估算一个集合中不同元素的数量。每个房间的 UV 统计只需要几十个字节但却能支撑亿级数据量下的基数统计。虽然它是概率估算误差大概在 0.81% 以内对于直播间的“在线人次”展示来说完全够用。第二层是用 Redis 的计数器来记录房间内当前在线人数在线连接数。这个数据相对更准确一些因为每次连接建立或断开时都会触发INCR或DECR操作public void onUserJoin(String roomId, String userId) { String key live:room:online: roomId; Long online redisTemplate.opsForValue().increment(key); redisTemplate.expire(key, 24, TimeUnit.HOURS); // 再异步同步一份到 MySQL用于对账和后续数据分析 } public void onUserLeave(String roomId, String userId) { String key live:room:online: roomId; redisTemplate.opsForValue().decrement(key); }这里有个细节是在线人数千万不能用 MySQL 的 COUNT 来实时算。我一开始用的就是SELECT COUNT(*) FROM user_online WHERE room_id?压测时发现每秒钟几百次这类查询就让数据库 CPU 飙到 100%。后来改成 Redis 计数性能直接提升了一个数量级。当然Redis 计数也有一个问题如果服务端异常重启计数可能会不准。所以我会用定时任务每隔几分钟把 Redis 里的在线人数快照同步到 MySQL用于数据对账。这样即使 Redis 数据丢了也能从 MySQL 恢复近似值。5. 常见问题与排查技巧实录5.1 问题速查表这套系统搭建过程中踩了不少坑我把最常见的问题整理成一个速查表方便你直接对照排查现象可能原因解决办法WebSocket 连接建立后几秒内被断开Nginx 未配置 Upgrade 转发或 read timeout 过短确认proxy_set_header Upgrade $http_upgrade把proxy_read_timeout调到 3600s客户端发消息后服务端没反应没有经过HttpObjectAggregator导致握手数据不完整在 Pipeline 里加入HttpObjectAggregator(65536)确保完整读取 HTTP 请求广播一条消息时所有房间的人都收到了用了全局 ChannelGroup 而不是按房间维度用ConcurrentHashMapString, ChannelGroup按 roomId 区分连接组压测时 CPU 飙升但连接数不高Netty 的 IO 线程里执行了阻塞操作如 Redis 同步调用把耗时操作改为异步线程池或消息队列服务端内存一直缓慢上涨每个连接创建了独立的定时任务或写入缓冲区堆积用全局调度任务代替每个连接的定时器检查 Channel 是否可写消息丢失消费者处理失败后没有重试或确认开启手动 ACK处理逻辑放在 try-catch 中失败后重试或进入死信队列客户端断线重连后收不到消息旧连接没有从 ChannelGroup 中移除覆盖了新的 Channel在channelInactive中彻底清理旧的 Channel确保 group 中和客户端的连接状态一致5.2 印象最深的三个坑第一个坑是文件描述符限制。我在本地 Mac 上做压测时连接数刚到 4000 多就报Too many open files。一开始以为是服务端出问题了排查半天才发现是客户端所在的机器我自己电脑先扛不住了。macOS 默认的ulimit -n是 256Linux 服务器上通常也是 1024。这个限制直接决定了一个进程能打开多少个文件描述符而每个 TCP 连接都要占用一个文件描述符。解决办法是临时调大ulimit -n 20000但要注意这个ulimit只在当前 shell 会话中生效如果通过 systemd 管理服务还需要在 service 文件里单独配置LimitNOFILE20000。我后来把这个写法也写进了部署文档不然换个机器又得重新踩一遍。第二个坑是ChannelGroup 并发移除导致的消息发向旧连接。直播间用户断线重连是很常见的高频场景如果旧的 Channel 没有及时清理重连成功后会出现“用户已经在新连接上但消息还在往旧连接发”的情况导致消息丢失。我在channelInactive里加了房间 ChannelGroup 的移除逻辑同时判断 Channel 是否 active 和 writable 后才发送才算治好了这个问题。第三个坑是Netty 的 write 不等于 flush。Netty 里channel.write()只是把消息放进缓冲区channel.flush()才真正触发网络发送。如果只是write而忘记flush短连接看不出来问题长连接场景下消息会一直堆积在缓冲区内存一路飙升最终可能把服务端拖死。我当时写广播逻辑时只调了channel.write()压测半小时后发现内存暴涨排查了好久才发现是这个原因。5.3 调优思路先找瓶颈再谈优化这次优化让我体会最深的一点是不要在没定位瓶颈之前盲目调参。很多网上教程上来就让你改 JVM 堆大小、改 GC 参数、改 Nginx worker 数量但如果你连瓶颈在哪儿都不知道这些操作就是碰运气。我的调优顺序是先做好监控至少要有 CPU、内存、GC、连接数这四类指标的可视化。压测过程中盯住指标变化看哪个指标先被打满。如果 CPU 先到瓶颈用jstack抓线程栈看是 IO 线程在忙还是业务线程在忙。如果内存先到瓶颈用jmap导出堆快照分析是连接对象占用多还是消息队列积压多。定位到瓶颈之后再动手优化每次只改一个参数改完重新压测对比数据。这套方法论看起来简单但特别管用。比如我发现内存高用jmap分析后发现大部分内存其实是 Netty 的堆外内存立刻就往调整写入缓冲区水位线和心跳机制的方向去查而不是盲目调 JVM 堆。方向对了优化的效率能翻好几倍。6. 这套环境的下一步扩展单机跑通后我心里清楚它距离“生产可用”还有一段距离但这段距离现在是可以踩着石头过河的了。我在笔记本上列了一个后续扩展清单列给你参考第一多节点部署与横向扩容。现在单机 Netty 能扛 1 万连接如果你要扛 10 万连接最简单的思路不是换更贵的服务器而是加机器。这时候就需要在前面那层 Nginx 的 upstream 里添加多台后端服务器让 Nginx 做连接级负载均衡。但连接分散之后会带来一个新问题用户 A 在节点 1 发了一条弹幕用户 B 在节点 2 怎么收到这就轮到 Redis Pub/Sub 出场了——连接层各节点订阅同一个 Redis 频道任何一条弹幕广播之后Redis 都会把消息推送给所有订阅了该频道的节点再由每个节点广播给自己维护的连接。这就是最经典的水平扩展方案。第二鉴权体系完善。当前版本里 WebSocket 握手时只校验了 URL 上带的 userId这个在生产环境完全不够。真实场景需要客户端先调用 REST 接口换取一个短期有效的 token然后 WebSocket 握手时带上 token服务端校验通过后才允许建立连接。同时还要做用户禁言、踢人下线等管理能力这些在直播场景里几乎是必需品。第三直播流媒体接入。这个项目目前只实现了“消息互动”这一层真正的直播视频流还完全没有接入。流媒体的方案通常是 RTMP 推流、HLS 拉流、WebRTC 低延迟传输这属于另一个技术栈camera、encoder、CDN的范畴和消息互动后端可以解耦。等消息层稳定了可以再单独学习这部分。第四链路追踪和可观测性。生产环境不能只靠日志定位问题特别是消息经过客户端、Nginx、Netty、MQ、MySQL 这条链路任何一个环节出问题都需要快速定位。可以用 SkyWalking 或者 Micrometer Prometheus 来做全链路监控把 JVM、连接数、队列积压情况、消息延迟全部暴露成指标接入告警。这一步越早做越好别等出了事故再补。最后分享一点个人心得整个过程走下来我最想分享给同样站在后端入门路口的朋友的一句话是别被“高并发”三个字吓住先定一个自己能验证的小目标把链路跑通再一步一步去压、去拆、去优化。这套环境从动手到跑通我用了大概两天时间真正值钱的部分其实不是那几行代码而是压测之后那些“现场排查”的经验。比如我永远都会记得第一次压到 1 万连接时看到 Netty 的 EventLoop 线程调度有条不紊、消息广播毫秒级送达时的那种兴奋感——那一刻我意识到后端高并发并没有想象中那么玄乎它不过是一套有规律可循的系统工程只要你肯从 0 开始一点一点把它搭起来。另外再透露一个小技巧当年我看源码和文档总觉得枯燥后来试了一个方法就是先给自己设定一个必须解决的问题再带着问题去查资料。就像这次做直播后端我先定了“10000 连接、1 秒延迟”的目标然后所有学习都围绕这个目标展开效率比我以前漫无目的地刷教程高了好几倍。这个思路你可以试试。
返回列表