我本身是做后端业务系统出身,日常接触最多的就是 CRUD 和接口联调,对高并发一直属于“听过没碰过”的状态。直到今年接到一个互动直播后台的活,要求支持万人同时在线、弹幕不卡、消息不乱,我整个人是从头皮麻到脚底。硬着头皮从零开始查资料、搭环境、压测、调优,折腾了将近三周,终于把这套环境给“手搓”出来了。这篇笔记就是完整记录我是怎么从一个不懂高并发的后端小白,一步步把直播涉及的长连接、消息削峰、状态同步、广播推送这些环节跑通的。整个过程算不上极致架构,但足够真实,也足够给和我一样想入门的后端同学参考。
1. 先搞明白:直播高并发到底在并发什么
学高并发最忌讳一上来就分布式、微服务、消息队列三板斧全招呼上。我第一周就是在这种“乱拳打死老师傅”的状态里度过的,最后发现连问题都没定义清楚。等我把直播场景拆开之后,才发现“高并发”这个词对不同模块的含义完全不一样。
1.1 直播系统的三条链路,压力完全不一样
一场直播从技术上拆,至少包含三条独立链路。
- 推流链路:主播手机或 OBS 把视频流推到服务器,这一路是单路到多路的复制,压力集中在带宽和转码,不在后端业务逻辑。
- 拉流链路:观众端从服务器拉取视频流,成千上万观众的观看压力主要由 CDN、边缘节点这一层承运,自建服务时压力也主要体现在带宽。
- 信令链路:弹幕、礼物、点赞、上下麦、在线人数统计这类实时交互,这才是后端高并发真正的主战场。
我把三条链路放进一个客厅场景类比一下:推流是嘉宾上台说话,拉流是观众坐在台下看,信令链路则是全场观众随时举手提问、交头接耳,主持人(后端)要把这些互动有序地调度起来。你可以想象,一群人的客厅里最累的不是音响,而是那个要回应所有人举手和递纸条的主持人。
所以真正考验后端高并发能力的,不是视频流本身,而是信令链路。这直接决定了我的资源投放方向和技术选型策略。
1.2 用数据量化:直播场景的压力值到底长什么样
理论讲完还不能落地,我需要把“万人直播”拆成具体数字。假设某直播间有 1 万观众同时在线,弹幕发送率高峰期约每秒 5% 的活跃用户在发消息,那就是每秒 500 条弹幕消息。再加上点赞、进场通知、礼物特效,瞬时消息量峰值冲到每秒 1500-2500 条非常正常。如果这些消息全部由后端实时广播给 1 万客户端,那么后端一秒要完成 1500 条消息的接收处理,再乘以 10000 个目标客户端的消息推送,理想情况下每秒触达次数就是 1500 万次。
这个数字一眼看上去很吓人,但它告诉我们一个道理:高并发直播架构的核心作用,其实是在“尽量少计算、尽量少复制、尽量少等待”的前提下,把一条消息推给所有人。说人话就是,后端要做的是一个缩小版的多对多消息系统,而不是一个增大版的 CRUD 管理系统。理解到这一层,后续的技术选型方向就清楚了:不能只用 Spring Boot 管好接口就收工,需要引入长连接管理、消息队列、广播网关等组件。
2. 架构设计:小白的第一个分层方案
脚手架搭之前,我拿了一整晚画系统的结构思路,画完又推倒重来三次。第一次画成了“全家桶”式微服务(网关、认证、用户、消息、房间、统计各一个服务),部署起来光启动就要五分钟,不适合个人从零到一。第二次妥协为单体 + 队列。最终定稿的是一个“既能跑得起来,又给未来留了拆分空间”的折中方案。
2.1 为什么后端必须“无状态化”和“横向扩展”
单人单机的后端应用,撑死能支撑几千个长连接。想要支撑数万人,唯一可靠的手段是加机器,这就要求应用本身是“无状态”的。怎么理解无状态?最直白的说法:任何一台后端服务器,都不知道某个用户上一次连的是哪一台。
举个现实场景:用户 A 第一次连上了服务器 X,弹幕也发得挺好。如果 A 断线重连时被负载均衡分到了服务器 Y,而服务器 Y 说“我不认识你”,那么用户体验就是掉线后离开直播间。如果服务器是“有状态”的(把用户信息、房间信息都放在本机内存里),这个情况就必然发生。
解决办法是让状态脱离单机进程,挪到集中式存储里,比如 Redis。后端只负责接收请求、向 Redis 读写数据,自身不保存关键业务状态。这样一来,任意后端服务器都能处理任意用户的请求,加机器就是一会的事。理解了“无状态化”,后面的连接池估算和会话管理思路就顺理成章了。
2.2 我的技术选型清单
踩过一次“啥都想上”的坑后,我最终的选型是这样的,给同样手搓环境的人一个参考。
| 层级 | 选型 | 说明 |
|---|---|---|
| 接入层 | Nginx | 反向代理、负载均衡、WebSocket 升级转发 |
| 业务 API | Spring Boot | 负责登录、房间信息、历史消息等常规接口 |
| 长连接网关 | Netty + WebSocket | 管理海量连接、心跳维护、消息推送 |
| 消息削峰 | RabbitMQ | 处理弹幕消息的写入与分发,避免高峰打垮数据库 |
| 状态缓存 | Redis | 存在线用户、房间人数、分布式会话标记 |
| 流媒体源站 | SRS | 收 RTMP 推流,输出 HTTP-FLV 给播放端 |
| 存储 | MySQL | 存弹幕记录、礼物记录、直播回看元数据 |
这套组合在不上 K8s 的前提下,用一台 8C16G 的服务器做压力验证,实测撑住了 3 万长连接加峰值 2000 QPS 的信令消息流。当然,这是压测环境的几何数,真实业务还得留出至少 3 倍冗余。
2.3 预估并发量:连接数、带宽、内存的计算方式
选型确定后,就要算资源。这步很实在,直接关系到要买多少钱的服务器,也关系到系统能不能扛住压力。
- 连接数估算:一个长连接在 Netty 中约占几十 KB 内存(包括 TCP 缓冲区、Channel 对象、业务扩展属性)。最粗略的估算公式为:
内存 = 预估连接数 × 单连接内存占用 × 冗余系数。按 2 万连接算,大约需要 2 万 × 50KB × 1.5 ≈ 1.5GB。留出操作系统和其他进程开销,8GB 内存的机器跑长连接网关是够的。 - 带宽估算:直播拉流带宽按码率算。1 路 1080P 直播码率差不多 4Mbps,1000 个观众同时观看就是 4Gbps 的下行流量,这个数字几乎不可能靠自建带宽硬扛,必须借助 CDN 的边缘分发能力。自建环境主要关注推流和测试拉流的带宽。
- 消息量估算:按一秒 1500 条信令消息、每条 1KB 计算,后端一秒钟需要消化 1.5MB 的数据。对于队列和 Redis 来说,这个量级并不高,真正的压力在全网广播时的扇出放大效应。
回看这套估算结果,我对“高并发直播”的核心诉求清晰了不少:限流靠 Nginx,连接靠 Netty,广播靠队列和发布订阅,状态靠 Redis,流媒体靠 SRS,剩下的才是业务逻辑。这个分工明确之后,就从“学习模式”切换成“建造模式”了。
3. 实操手记:从零搭出完整环境
这一部分直接上实操。我尽量按时间顺序记录,方便想复现的朋友跟着操作。
3.1 第一步:搭建流媒体服务,让推拉流先通起来
直播环境里最先要解决的是“流怎么走通”的问题。视频流处理不能靠 Java 硬写,直接上开源方案最靠谱。我这里选用了 SRS(Simple Realtime Server),它支持 RTMP 接入,又能输出 HTTP-FLV 给网页播放器,非常契合“主播推流,用户浏览器拉流”的场景。
我用 Docker 方式启动 SRS,配置文件简单到让人觉得是在做填空题:
docker run -d --name srs \ -p 1935:1935 -p 8080:8080 -p 1985:1985 \ -v /home/srs/conf/srs.conf:/usr/local/srs/conf/srs.conf \ registry.cn-hangzhou.aliyuncs.com/ossrs/srs:5.0srs.conf 里最核心的配置就两段。
listen 1935; max_connections 10000; http_api { enabled on; listen 1985; } http_server { enabled on; listen 8080; dir ./objs/nginx/html; } vhost __defaultVhost__ { http_remux { enabled on; mount [vhost]/[app]/[stream].flv; hstrs on; } }配置完成后,我本地用 OBS 推流,推流地址填入:
rtmp://服务器IP:1935/live/test播放端用 HTTP-FLV 地址验证:
http://服务器IP:8080/live/test.flv用 VLC 打开这个地址,看到画面动起来的那一刻,我整个人松了一口气。流媒体通,后面的一大半问题就不是流本身,而是消息和状态的并发了。
一个小提醒:HTTP-FLV 兼容性极佳,适合快速验证。如果以后要做连麦或者更低的延迟,那就要考虑 WebRTC 和 SRS 的 RTC 配置,这属于进阶话题,等第一版跑通后再折腾不迟。
3.2 第二步:搭建 Netty 网关,管理上万条长连接
视频流通了,接下来是最硬核的环节:长连接网关。为什么不用 Spring Boot 自带的 WebSocket 而要用 Netty?原因有两个:一是 Netty 对海量连接的内存管理更精细,二是它可以做到极致的 IO 线程模型,避免业务线程阻塞。
网关的核心职责是:接收客户端的 WebSocket 升级请求,鉴权通过后建立连接,维护连接与房间、用户的关系,处理心跳,把服务端收到的消息转交给 MQ。
连接分组管理是重点。我设了一个ChannelManager来维护两种映射:
userId -> Channel:用于给指定用户推送消息(比如私信、礼物回调)。roomId -> ChannelGroup:用于给整个房间广播(比如弹幕、系统消息)。
核心代码结构大概是这样:
public class ChannelManager { // userId -> Channel private static final ConcurrentHashMap<Long, Channel> USER_CHANNEL = new ConcurrentHashMap<>(); // roomId -> ChannelGroup private static final ConcurrentHashMap<Long, DefaultChannelGroup> ROOM_GROUP = new ConcurrentHashMap<>(); public static void addUserToRoom(Long userId, Long roomId, Channel channel) { USER_CHANNEL.put(userId, channel); DefaultChannelGroup group = ROOM_GROUP.computeIfAbsent(roomId, k -> new DefaultChannelGroup(GlobalEventExecutor.INSTANCE)); group.add(channel); channel.attr(AttributeKey.valueOf("roomId")).set(roomId); } public static void removeUser(Long userId, Channel channel) { USER_CHANNEL.remove(userId); Long roomId = channel.attr(AttributeKey.valueOf("roomId")).get(); if (roomId != null) { DefaultChannelGroup group = ROOM_GROUP.get(roomId); if (group != null) { group.remove(channel); if (group.isEmpty()) { ROOM_GROUP.remove(roomId); } } } channel.close(); } public static void broadcast(Long roomId, Object msg) { DefaultChannelGroup group = ROOM_GROUP.get(roomId); if (group != null) { group.writeAndFlush(new TextWebSocketFrame(JSON.toJSONString(msg))); } } }这里有个性能要点:广播时使用DefaultChannelGroup的底层批量写能力,它会遍历组内 Channel 并调用writeAndFlush,避免自己写 for 循环带来的上下文切换开销。
建连后,客户端和服务端需要保持心跳。我的方案是客户端每 30 秒发送一个 ping 帧,服务端超过 90 秒没收到就主动断开。Netty 的IdleStateHandler可以直接支持这个逻辑:
ch.pipeline().addLast(new IdleStateHandler(90, 30, 0, TimeUnit.SECONDS));服务器 90 秒没收到客户端数据就触发关闭,客户端 30 秒没发送数据就触发一个用户自定义事件(可用来探测心跳失败并重拾连接)。这个“30 秒 PING、90 秒超时”参数组合是我调了几次之后的稳定方案,太短会误杀正常用户(比如切后台时网络休眠),太长又会让脏连接占着内存。
3.3 第三步:引入 RabbitMQ,把突增流量削平
弹幕消息如果直接写进数据库或直接广播,瞬间的高峰流量会直接打满 DB 连接和网络 I/O。这里必须有队列削峰。
我的做法是:
- 客户端把消息发到网关,网关不做业务处理,直接投递到 RabbitMQ 的
danmu.queue。 - 独立的消费者服务从队列里拉取消息,做敏感词过滤、存入 Redis 和 MySQL,再调用网关的广播接口推送给房间内所有客户端。
为什么这么绕?核心原因是削峰填谷。假设高峰期一秒涌入 5000 条弹幕,如果直接让业务服务处理,MySQL 瞬间就挂了。队列把 5000 条消息按消费者的处理能力(比如每秒 1000 条)平稳消费掉,高峰期多出来的部分先囤在队列里,低峰期再慢慢消费完。用户侧感受到的推送延迟差个几百毫秒,完全可以接受。
RabbitMQ 的核心配置如下:
@Configuration public class RabbitConfig { @Bean public Queue danmuQueue() { Map<String, Object> args = new HashMap<>(); // 队列最大长度,避免消息无限堆积 args.put("x-max-length", 100000); // 超出长度的消息进入死信队列,便于后续排查 args.put("x-dead-letter-exchange", "dlx.exchange"); args.put("x-dead-letter-routing-key", "danmu.dlx"); return new Queue("danmu.queue", true, false, false, args); } @Bean public DirectExchange dlxExchange() { return new DirectExchange("dlx.exchange"); } @Bean public Queue dlxQueue() { return new Queue("danmu.dlx.queue", true); } @Bean public Binding dlxBinding() { return BindingBuilder.bind(dlxQueue()).to(dlxExchange()).with("danmu.dlx"); } }3.4 第四步:用 Redis 保存实时状态,避免单机内存失效
在一个多节点环境中,Netty 网关的所有状态都不能只留在本机内存,否则断线重连到另一台机器就直接“失忆”。这一步用 Redis 集中管理状态:
- 在线用户标记:
SET online:userId roomId EX 120,只要用户心跳就续期。这样用户是否在直播间独立于单台服务器判断。 - 房间人数计数:
INCR room:count:{roomId},配合DECR处理退出,能实时统计在线人数。 - 分布式会话标记:
SET session:userId gatewayNodeId EX 120,这样未来做“定向推送”(用户在指定网关节点)时有迹可循。
代码层面对 Redis 的操作建议用连接池,别每次操作都新建连接,实测能够显著降低延迟:
@Bean public LettuceConnectionFactory redisConnectionFactory() { RedisStandaloneConfiguration config = new RedisStandaloneConfiguration(); config.setHostName("localhost"); config.setPort(6379); config.setDatabase(0); return new LettuceConnectionFactory(config); } @Bean public RedisTemplate<String, String> redisTemplate() { RedisTemplate<String, String> template = new RedisTemplate<>(); template.setConnectionFactory(redisConnectionFactory()); template.setKeySerializer(new StringRedisSerializer()); template.setValueSerializer(new Jackson2JsonRedisSerializer<>(Object.class)); return template; }有一个很关键的细节:Redis 的keys命令在生产环境千万别用,一旦在线人数上万,keys online:*会直接卡死 Redis。要么用SCAN游标遍历,要么直接按房间维度拆分 key 结构,比如room:10001:onlineSet存一个房间的在线用户 ID 集合。
4. 压测与调优:确认这套环境真的能扛
搭建完成后,我没有直接上线,而是进行了为期两天的压测与调优。如果你不想在观众面前出丑,这一步绝对省不得。
4.1 压测工具与脚本:模拟万人在线
压测 WebSocket 长连接没有现成点击工具,这里我用了一段简单的 Java 压测客户端,核心原理是循环创建 WebSocket 连接,并定时发送弹幕消息。
public class PressureClient { public static void main(String[] args) throws Exception { int totalConnections = 10000; CountDownLatch latch = new CountDownLatch(totalConnections); ExecutorService executor = Executors.newFixedThreadPool(500); for (int i = 0; i < totalConnections; i++) { final int userId = i; executor.submit(() -> { try { WebSocketClient client = new WebSocketClient(new URI("ws://127.0.0.1:8080/ws?userId=" + userId)) { @Override public void onOpen(ServerHandshake handshakedata) { latch.countDown(); } }; client.connect(); // 每隔10秒发一条弹幕 if (userId % 10 == 0) { client.send("hello from user " + userId); } } catch (Exception e) { e.printStackTrace(); } }); } latch.await(30, TimeUnit.SECONDS); System.out.println("连接完成"); } }实测下来,这个压测客户端能同时撑起 1 万连接,也基本够用。专业一点的压测工具可以考虑 Gatling 或 JMeter 的 WebSocket Sampler,但学习成本会更高。
4.2 遇到并解决的三个瓶颈
第一版压测结果非常惨烈:连接数超过 3000 就大量超时,弹幕发送后客户端普遍延迟 2 秒以上。我逐一排查,找出三个致命瓶颈。
第一个是文件描述符限制。Linux 默认单进程最大文件描述符是 1024,也就是一个进程最多能同时打开 1024 个 socket。这个不调,一万连接无从谈起。解决方式:
ulimit -n 1000000为了让重启后依然生效,需要写入/etc/security/limits.conf:
* soft nofile 1000000 * hard nofile 1000000第二个是 Nginx 的默认超时。Nginx 默认读取超时只有 60 秒,WebSocket 长连接如果 60 秒没有数据交互,连接就被切断。心跳虽然存在,但心跳间隔和 Nginx 超时之间必须留足余量。我最终把 Nginx 配置调整成这样:
server { listen 80; server_name live.example.com; location /ws { proxy_pass http://127.0.0.1:8080; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection "upgrade"; proxy_set_header Host $host; proxy_set_header X-Real-IP $remote_addr; proxy_read_timeout 120s; proxy_send_timeout 120s; } }第三个是 Netty 的 IO 线程数与业务线程池配置。默认的 Boss/Worker 线程数对高连接数并不友好,我按 CPU 核数调整如下:
EventLoopGroup bossGroup = new NioEventLoopGroup(1); EventLoopGroup workerGroup = new NioEventLoopGroup(Runtime.getRuntime().availableProcessors() * 2);同时,业务处理(比如 JSON 解析、Redis 写入)不能直接放到 IO 线程,必须dispatch到独立的业务线程池,否则 IO 线程会被慢业务拖死。
4.3 压测数据对比:调优是真实有效果的
调优前后对比差异明显,我记录了一组压测数据供参考(测试环境:8C16G 云服务器,单节点 Netty 网关,RabbitMQ 消费服务独立部署)。
| 指标 | 调优前 | 调优后 |
|---|---|---|
| 最大稳定连接数 | 3000 | 30000 |
| 弹幕消息峰值接收(QPS) | 300 | 2100 |
| P99 弹幕广播延迟 | 2.3秒 | 110毫秒 |
| CPU 使用率(峰值) | 85% | 70% |
| 内存使用率 | 60% | 52% |
最令人惊喜的是 P99 延迟下降了近 20 倍,从无法用的状态变成了流畅体验。这个结果让我确信:瓶颈不在服务器本身,而在配置和代码细节的优化。
5. 常见问题与排查技巧实录
学高并发一个很大的绊脚石是出了问题不知道怎么排查。这一节我把自己遇到的典型问题和分析方式整理出来,也算是一份避坑清单。
5.1 连接不上?先分清是哪一层的锅
排查网络问题有一个铁律:别猜,分好层,逐层验证。
我遇到过客户端一直报 WebSocket 握手失败的问题。先 curl 一下 Nginx 端口通不通,发现正常;再用 WebSocket 客户端直连后端端口,也正常;问题锁定在 Nginx 转发层。打开 Nginx error log 才发现/ws路径没有匹配到代理规则,请求被当成静态文件处理。
还有一次最隐蔽:服务器安全组防火墙只放行了 80 端口,WebSocket 使用的 8080 端口没放行。外部无法连接,本地测试却一切正常。这类问题一半时间都浪费在症状排查上,建议一开始就用telnet 服务器IP 端口做一个端口连通性测试,能省半天。
5.2 消息延迟高?不要只盯消费者
弹幕从发送到接收,经历了客户端 -> 网关 -> MQ -> 消费 -> 广播 -> 客户端六个环节。延迟高时只盯一个环节很容易漏。
排查重点应该放在三处:
- 生产端:网关投递消息到 MQ 的速度是否稳定。如果网关的
channelRead业务线程池满了,消息投递会排队。 - MQ 消费:查看队列堆积情况,如果消费速率小于生产速率,队列长度会持续增长,延迟自然上升。
- 广播端:Netty 广播的 Group 大小和写缓冲水位线设置是否合理,如果写出速度被对端慢连接拖累,会影响整个房间的推送效果。
我的一个重点教训是:一台慢客户端会拖垮一个房间的广播。后来在ChannelGroup的写出缓冲上设置了WRITE_BUFFER_WATER_MARK,并将慢连接自动降级或断开,问题才得到缓解。
5.3 长连接数上不去?不一定是后端的问题
连接数上不去,80% 的原因不是 Java 代码,而是操作系统参数。除了上面提到的 ulimit,还有两个关键参数值得检查:
# 查看当前连接数 netstat -an | grep ESTABLISHED | wc -l # 临时调整 TCP 连接复用能力 sysctl -w net.ipv4.tcp_tw_reuse=1 sysctl -w net.ipv4.ip_local_port_range="1024 65535"通常压测环境最容易踩的坑是客户端端口不够用。作为客户端发起大量短连接时,源端口迅速耗尽,表现为“后面几百个连接全失败”。调整ip_local_port_range往往立竿见影。
5.4 敏感词过滤:不要小瞧它的性能消耗
弹幕消息经过 RabbitMQ 后,消费者内部要做敏感词匹配。如果每次用 Java 正则表达式扫描,一条弹幕几千次匹配,高峰期 CPU 直接拉满。
我的优化思路是:构建一个基于前缀树的敏感词过滤工具,把敏感词库加载到内存,扫描时逐字符走 DFA(确定性有限自动机),完全避开正则表达式的回溯开销。实测单条弹幕从平均 5ms 降到 0.2ms,消费者处理能力提升了 20 倍。
这个细节说明,高并发优化往往不是靠某个宏大架构,而是靠无数个类似的前端小优化积累出来的。
写在最后的一点心里话
这次从零搭建直播高并发环境,让我最大的收获其实不是学会了某个工具,而是搞清楚了一条思路:高并发的本质不是“机器多”,而是“少干活”。把每次请求都尽量做得轻,把不紧急的活都丢进队列,把易变的状态都放进缓存,把难维护的连接都交给专门的网关层,这套组合拳下来,单机也能扛出一个不错的量级。说实话,第一次看到 3 万长连接同时在线的压测曲线时,我心里是有一种朴素的成就感的,那种感觉就像一个新手司机第一次独自开上了高速。这套方案还远谈不上完善,后续如果要做多地容灾、要做动态扩缩容,还有很长的路要走。但至少对我这样的后端小白来说,它证明了“高并发直播”没有传说中那么遥不可及,敢动手拆解和尝试,就已经迈过了最难的那道坎。