拓冰建站拓冰建站
首页 / 资讯中心 / 正文

从零搭建直播高并发环境:后端小白的实战笔记与架构设计

1. 从零搭建直播高并发环境一个后端小白的实战笔记直播这玩意儿看别人做觉得挺简单——推个流、拉个流、弹幕聊聊天能有多难真到自己动手搭一套能扛住高并发的环境才发现每一步都是坑。我是一名后端开发之前主要写业务接口对直播领域几乎零经验。这个项目是我从零开始一步步把一套直播高并发环境搭起来的完整记录。里面涉及架构设计、技术选型、Redis 缓存设计、高并发 IM、前后端分离等核心环节适合有一定后端基础、想了解直播系统怎么扛并发的朋友参考。我不会讲太多虚的主要说清楚每个决策背后的逻辑、踩过的坑以及可以直接抄作业的配置和代码。整个项目从最初的单机 demo到后来能支撑万级并发在线中间经历了三次比较大的架构调整。第一次是发现单节点扛不住弹幕写入第二次是消息延迟飙升第三次是缓存击穿导致数据库差点崩掉。每一次调整都让我对“高并发”这三个字有了更具体的理解。下面我把整个过程拆开来讲尽量把每个环节的“为什么”说透。2. 整体架构设计与技术选型思路2.1 为什么最终选了这套组合一开始我想得很简单Spring Boot 写几个接口前端 Vue 调一下消息用 WebSocket 推数据库 MySQL 存着完事。结果本地测试没问题一上压测工具500 并发就开始各种超时。问题出在哪直播场景和普通业务系统最大的区别在于写少读多、消息密集、实时性要求高。一场直播可能有几万人在线但真正发弹幕的可能只有几百人大部分人是在看。这意味着读请求远远大于写请求而且读请求对延迟极其敏感。基于这个特点我最终的技术选型是这样的组件选型核心原因后端框架Spring Boot NettyNetty 处理 WebSocket 长连接比 Tomcat 的 NIO 更可控消息队列RocketMQ弹幕削峰填谷保证消息不丢缓存Redis Cluster扛住高频读弹幕列表、在线状态都走缓存数据库MySQL 分表只存核心数据弹幕历史归档前端Vue3 WebSocket前后端分离实时消息用原生 WebSocket压测JMeter模拟高并发场景验证承载能力这里重点说下为什么用 Netty 而不是直接用 Spring WebSocket。Spring 的 WebSocket 底层也是 Netty但它封装得太厚连接数一上去内存和 GC 都不好控制。自己用 Netty 搭可以精细控制 ByteBuf 的分配和释放对长连接场景更友好。当然代价是代码复杂度上去了这个后面会讲。2.2 架构分层与数据流向整个系统的数据流向是这样的用户通过前端发起 WebSocket 连接Netty 网关接收后把消息投递到 RocketMQ消费端处理后再写 Redis 和 MySQL同时通过 Netty 推送给房间内其他用户。读请求直接走 Redis不碰数据库。这个设计的关键在于把写路径和读路径彻底分开。写路径经过 MQ 异步化保证即使数据库慢也不会阻塞消息推送读路径全部走缓存保证毫秒级响应。听起来简单但实际落地时MQ 的消费速度、Redis 的容量规划、Netty 的连接管理每一个都是坑。提示不要一上来就搞微服务。我最初就是被“分布式架构设计”这个词带偏了拆了五六个服务结果运维成本极高问题排查困难。后来合并成三个核心服务网关服务、消息服务、业务服务反而更稳。2.3 容量预估与资源规划做高并发系统拍脑袋定配置是大忌。我当时的预估逻辑是这样的假设峰值在线 5 万人其中 10% 的人会发弹幕也就是 5000 人同时发消息。每人平均 10 秒发一条那么 QPS 大约是 500。但直播场景有突发性比如主播抽奖瞬间弹幕量可能翻 10 倍所以按 5000 QPS 来设计。Redis 这边每个房间的弹幕列表保留最近 200 条每条弹幕按 200 字节算一个房间 40KB。1000 个房间就是 40MB完全放得下。在线状态用 Hash 存储每个用户 50 字节5 万用户也就 2.5MB。所以 Redis 的内存压力不大主要压力在网络 IO 上。MySQL 这边弹幕历史按天分表每天一张表保留 30 天。单表数据量控制在 500 万以内查询走索引问题不大。3. 核心细节解析与实操要点3.1 Netty 网关的连接管理Netty 网关是整个系统的入口负责维护所有 WebSocket 连接。这里最大的坑是连接数上不去和内存泄漏。我一开始用默认配置单机只能撑 3000 连接左右后来调整了以下几个参数// Netty 服务端配置 EventLoopGroup bossGroup new NioEventLoopGroup(1); EventLoopGroup workerGroup new NioEventLoopGroup(16); // 根据 CPU 核数调整 ServerBootstrap b new ServerBootstrap(); b.group(bossGroup, workerGroup) .channel(NioServerSocketChannel.class) .option(ChannelOption.SO_BACKLOG, 1024) .childOption(ChannelOption.SO_KEEPALIVE, true) .childOption(ChannelOption.TCP_NODELAY, true) .childOption(ChannelOption.SO_RCVBUF, 32 * 1024) .childOption(ChannelOption.SO_SNDBUF, 32 * 1024) .childHandler(new ChannelInitializerSocketChannel() { Override protected void initChannel(SocketChannel ch) { ch.pipeline().addLast(new HttpServerCodec()); ch.pipeline().addLast(new HttpObjectAggregator(65536)); ch.pipeline().addLast(new WebSocketServerProtocolHandler(/ws)); ch.pipeline().addLast(new LiveWebSocketHandler()); } });SO_BACKLOG调到 1024 是为了应对突发连接TCP_NODELAY关闭 Nagle 算法保证消息实时性。SO_RCVBUF和SO_SNDBUF调小是为了降低内存占用因为直播消息都很小不需要大缓冲区。连接管理方面我用了一个ConcurrentHashMap来保存userId - Channel的映射同时用另一个 Map 保存roomId - SetuserId。这样推送消息时先根据 roomId 找到所有 userId再找到对应的 Channel 推送。注意这里一定要用ConcurrentHashMap否则并发场景下会出问题。注意Channel 的writeAndFlush是异步的如果客户端网络慢消息会堆积在 ChannelOutboundBuffer 里导致内存暴涨。我的做法是加一个水位监听超过阈值就断开连接或者丢弃消息。3.2 Redis 缓存设计与高并发读写Redis 在这个系统里承担了三个角色弹幕缓存、在线状态、分布式锁。弹幕缓存用 List 结构每个房间一个 key比如live:danmaku:room:1001用LPUSH写入LTRIM保留最近 200 条。读的时候用LRANGE一次性取出来。在线状态用 Hashkey 是live:online:room:1001field 是 userIdvalue 是时间戳。用户进入房间时HSET离开时HDEL。统计在线人数用HLENO(1) 复杂度非常快。分布式锁用在“房间创建”和“用户禁言”这两个场景。用SET key value NX PX 3000实现value 用 UUID释放锁时用 Lua 脚本保证原子性。if redis.call(get, KEYS[1]) ARGV[1] then return redis.call(del, KEYS[1]) else return 0 end这里有个坑缓存击穿。热门房间的弹幕列表如果过期了大量请求同时打到数据库数据库瞬间就挂了。我的解决方案是弹幕列表永不过期但设置一个逻辑过期时间后台定时任务刷新。这样既保证了数据新鲜度又避免了击穿。另一个坑是大 key。一个房间的弹幕列表如果无限增长会变成大 key操作时阻塞 Redis。所以一定要用LTRIM限制长度200 条足够了用户也不会翻太久。3.3 RocketMQ 削峰填谷与消息可靠性弹幕消息不直接写库而是先投递到 RocketMQ消费端再慢慢处理。这样做的好处是即使数据库压力大消息也不会丢只是延迟增加。RocketMQ 的配置要点DefaultMQProducer producer new DefaultMQProducer(live_producer_group); producer.setNamesrvAddr(127.0.0.1:9876); producer.setRetryTimesWhenSendFailed(3); producer.setSendMsgTimeout(3000); producer.start(); Message msg new Message(live_topic, danmaku, JSON.toJSONBytes(danmaku)); SendResult result producer.send(msg);消费端用MessageListenerConcurrently并发消费线程数根据 CPU 核数调整。注意消费端一定要做幂等处理因为 RocketMQ 保证至少投递一次可能重复。我的做法是用弹幕 ID 做去重Redis 里SETNX一个 key过期时间 5 分钟。提示RocketMQ 的 Topic 要提前创建并且设置合理的队列数。队列数决定了并发消费的并行度我设的是 8 个队列对应 8 个消费线程。3.4 前后端分离下的 WebSocket 通信前端用 Vue3WebSocket 连接封装在useWebSocket这个 composable 里。核心逻辑是连接建立后监听onmessage收到消息后根据类型分发到不同的处理函数。心跳包每 30 秒发一次服务端如果 60 秒没收到心跳就断开连接。const ws new WebSocket(wss://live.example.com/ws?token token); ws.onopen () { setInterval(() { ws.send(JSON.stringify({ type: heartbeat })); }, 30000); }; ws.onmessage (event) { const data JSON.parse(event.data); if (data.type danmaku) { danmakuList.value.push(data.payload); } };跨域问题在 WebSocket 里不存在因为 WebSocket 不受同源策略限制。但鉴权要做在连接建立时通过 URL 参数传 token服务端校验通过后才允许连接。不要等到连接建立后再发鉴权消息那样会浪费连接资源。4. 实操过程与核心环节实现4.1 环境准备与基础服务搭建我用的是一台 8 核 16G 的云服务器操作系统 Ubuntu 22.04。基础服务包括 MySQL 8.0、Redis 7.0、RocketMQ 5.0全部用 Docker 部署方便管理。# Redis 集群部署3 主 3 从 docker run -d --name redis-node-1 --net host redis:7.0 \ redis-server --port 7001 --cluster-enabled yes \ --cluster-config-file nodes.conf --cluster-node-timeout 5000 \ --appendonly yesMySQL 的配置重点是max_connections调到 2000innodb_buffer_pool_size调到 8Ginnodb_flush_log_at_trx_commit设为 2牺牲一点持久性换性能直播场景可以接受。RocketMQ 的 Broker 配置里brokerRole设为ASYNC_MASTERflushDiskType设为ASYNC_FLUSH这样写入性能最好。当然如果对消息可靠性要求极高可以改成同步刷盘但吞吐量会下降。4.2 核心代码实现弹幕发送与推送弹幕发送的完整流程是这样的前端通过 WebSocket 发送弹幕内容Netty 网关收到后先做参数校验和敏感词过滤然后投递到 RocketMQ同时返回一个 ack 给前端。消费端从 MQ 拉取消息写入 Redis 和 MySQL然后通过 Netty 推送给房间内所有用户。// Netty 网关处理弹幕 protected void channelRead0(ChannelHandlerContext ctx, TextWebSocketFrame frame) { DanmakuRequest request JSON.parseObject(frame.text(), DanmakuRequest.class); // 参数校验 if (StringUtils.isBlank(request.getContent()) || request.getContent().length() 100) { ctx.writeAndFlush(new TextWebSocketFrame(参数错误)); return; } // 敏感词过滤 if (sensitiveFilter.contains(request.getContent())) { ctx.writeAndFlush(new TextWebSocketFrame(内容包含敏感词)); return; } // 投递到 MQ Message msg new Message(live_topic, danmaku, JSON.toJSONBytes(request)); producer.send(msg); // 返回 ack ctx.writeAndFlush(new TextWebSocketFrame(ok)); }消费端的逻辑consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) - { for (MessageExt msg : msgs) { DanmakuRequest request JSON.parseObject(msg.getBody(), DanmakuRequest.class); // 幂等去重 String key danmaku:dedup: request.getId(); Boolean success redisTemplate.opsForValue().setIfAbsent(key, 1, 5, TimeUnit.MINUTES); if (Boolean.FALSE.equals(success)) { continue; } // 写 Redis String redisKey live:danmaku:room: request.getRoomId(); redisTemplate.opsForList().leftPush(redisKey, JSON.toJSONString(request)); redisTemplate.opsForList().trim(redisKey, 0, 199); // 写 MySQL异步 danmakuMapper.insert(request); // 推送给房间内用户 pushService.pushToRoom(request.getRoomId(), request); } return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; });推送服务里从roomUserMap拿到所有 userId再遍历userChannelMap拿到 Channel调用writeAndFlush。注意这里要加异常处理如果 Channel 已经关闭要从 Map 里移除否则会内存泄漏。4.3 压测方案与性能调优压测用 JMeter模拟 5000 个并发用户每个用户建立 WebSocket 连接每秒发送一条弹幕。压测脚本里用WebSocket Sampler配置好连接地址和消息内容。第一次压测结果很惨5000 并发下消息延迟平均 2 秒P99 达到 8 秒。排查后发现三个问题一是 Netty 的 worker 线程数不够二是 RocketMQ 的队列数太少三是 Redis 的LPUSH和LTRIM不是原子操作导致数据不一致。调优措施worker 线程数从 8 调到 16RocketMQ 队列数从 4 调到 8Redis 操作用 Lua 脚本合并成原子操作。再次压测延迟降到 200msP99 降到 800ms基本达标。-- 原子写入弹幕并裁剪 redis.call(lpush, KEYS[1], ARGV[1]) redis.call(ltrim, KEYS[1], 0, 199) return 1注意压测时一定要监控服务器资源。我用top和iostat看 CPU 和磁盘 IO用redis-cli --latency看 Redis 延迟用 RocketMQ 的控制台看消息堆积情况。任何一个指标异常都要停下来分析。5. 常见问题与排查技巧实录5.1 连接数上不去怎么办这是最常见的问题。现象是压测到一定并发后新连接建立失败。排查思路先看ulimit -n默认是 1024改成 65535。再看 Netty 的SO_BACKLOG默认是 128改成 1024。最后看操作系统的somaxconnsysctl -w net.core.somaxconn1024。这三个地方都改了单机撑 1 万连接没问题。5.2 消息延迟突然飙升延迟飙升通常是 MQ 堆积或者 Redis 慢查询导致的。先看 RocketMQ 的消费进度如果堆积量大说明消费端处理不过来要么加消费线程要么优化消费逻辑。再看 Redis 的slowlog如果有LRANGE大 key 的操作说明弹幕列表太长了要检查LTRIM是否生效。5.3 缓存击穿导致数据库崩溃前面提过热门房间的弹幕列表如果过期大量请求会打到数据库。解决方案是逻辑过期 后台刷新。具体做法是Redis 里的弹幕列表不设过期时间但每个房间额外存一个expire_time后台任务定时检查快过期时主动刷新。这样用户永远读不到空数据数据库也不会被击穿。5.4 常见问题速查表问题现象可能原因排查方法解决方案连接建立失败文件描述符限制ulimit -n调到 65535消息延迟高MQ 堆积控制台看消费进度加消费线程或优化逻辑Redis 超时大 key 操作slowlog限制 key 长度数据库 CPU 高缓存击穿看 QPS 和缓存命中率逻辑过期 后台刷新内存泄漏Channel 未移除jmap看对象数量关闭时从 Map 移除5.5 独家避坑技巧第一个技巧Netty 的Channel一定要在channelInactive回调里从 Map 移除否则连接断开后 Map 越来越大最后 OOM。我一开始就是忘了这个跑了一天内存就满了。第二个技巧RocketMQ 的消息 ID 在重试时会变所以幂等去重不能用消息 ID要用业务 ID。我在弹幕请求里生成了一个 UUID 作为业务 ID消费端用这个去重。第三个技巧Redis 的LPUSH和LTRIM分开执行时如果中间服务挂了会导致列表无限增长。用 Lua 脚本合并成原子操作要么都成功要么都失败。第四个技巧压测时不要只压 WebSocket还要压 HTTP 接口。因为用户进入房间时会调 HTTP 接口获取历史弹幕这个接口如果慢用户会卡在加载页面。6. 后续扩展与个人体会这套环境目前能稳定支撑 1 万并发在线峰值 5000 QPS 的弹幕写入。后续如果要扩展到 10 万并发需要做几件事一是 Netty 网关做集群用 Nginx 做四层负载均衡二是 Redis 做集群分片按房间 ID 哈希三是 RocketMQ 加 Broker 节点提升吞吐量。我个人在实际操作中的体会是高并发系统不是靠堆配置堆出来的而是靠对业务场景的深刻理解。直播场景的核心是“实时”和“不丢消息”所有的技术选型和架构设计都要围绕这两个点。不要盲目追求新技术适合场景的才是最好的。比如我一开始想用 Kafka后来发现 RocketMQ 的延迟更低更适合直播场景。最后再分享一个小技巧做压测时先用小并发跑通流程再逐步加大并发。每次只调一个参数观察效果。这样出了问题容易定位不会一团乱麻。我见过很多人一上来就压 1 万并发结果系统崩了都不知道哪里崩的。
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门