🚀 拒绝“单机版”教学!Netty+Redis+RocketMQ 打造百万级实时互动引擎 | 智学云开发日记(五)

作者:StarLiuBB
标签:#Netty #Redis #高并发 #WebSocket #架构设计 #Java
背景:当在线课堂变成“单机游戏”,如何用技术挽回百万学生的体验?


👋 开场白:当直播课变成了“默剧”

“老师,我刚去上了个厕所,这题选C还是选D啊?”
“老师,屏幕糊得像马赛克,是我网卡了吗?”

还记得智学云旧版平台吗?那时候学生发出的呐喊就像扔进了黑洞。为什么?因为我们当初图省事,用的是传统的 HTTP 轮询(Polling)

学生想问问题,前端每隔几秒去戳一下服务器:“有新消息吗?”。这哪是直播课啊,简直是大型**“单机游戏”**现场,老师在台上自嗨,学生在台下挂机。

拒绝摆烂! 为了让课堂“燃”起来,我们决定引入实时互动服务。弹幕起飞、点赞刷屏、收藏秒存——我们要让在线课堂拥有比拟 B 站的丝滑体验!

今天,不开玩笑,全是干货。带大家复盘一下我们是如何用 Netty + Redis + RocketMQ 这套“三驾马车”,扛住流量洪峰,实现毫秒级互动的。


🧐 Round 1:HTTP vs WebSocket,谁才是YYDS?

项目启动时,团队内部发生了一场激烈的**“路线之争”**。

❌ 方案A:HTTP 轮询(那个“死缠烂打”的前任)

前端小哥提议:“简单点,我写个 setInterval,每3秒问一次服务器有没有新消息。”

// 前端疯狂试探:服务器你理理我啊!
setInterval(() => {
    fetch('/api/messages').then(data => updateUI(data));
}, 3000); 

后果预演:

  1. 无效请求爆炸:几千个学生同时在线,QPS瞬间破万,服务器CPU直接破防,而其中99%的请求都是空的(没有新消息)。
  2. 延迟感人:老师刚讲了个段子,学生3秒后才笑,这反射弧比树懒还长。
  3. 电量杀手:手机发烫,电量尿崩,学生直呼“退钱”。

✅ 方案B:WebSocket 长连接(双向奔赴的真爱)

经过技术委员会(其实就是我和架构师)拍板,必须上 WebSocket

  • 全双工通信:服务器想推就推,客户端想发就发,拒绝“单相思”。
  • 低开销:建立一次连接,终身受益,头部信息少,带宽省到家。
  • 实时性:毫秒级触达,老师说“扣1送分”,屏幕瞬间全是“111”。

结论: WebSocket 才是实时互动的版本答案


🏗️ Round 2:架构设计——如何优雅地扛住百万并发?

光有 WebSocket 协议不够,我们得有能承载它的容器。Tomcat?并发能力差点意思。最终我们选择了 Java 界的网络编程霸主 —— Netty

🗺️ 全景架构图(高逼格版)

CSDN支持Mermaid渲染,以下架构图会自动生成:

核心互动层

WebSocket

广播

广播

👶 学生/老师

⚖️ 负载均衡 Nginx

🔥 Netty Server 1

🔥 Netty Server 2

🔥 Netty Server N...

📣 Redis Pub/Sub

🚀 RocketMQ

💾 消息持久化服务

🐬 MySQL

🧩 核心组件职责(不讲虚的)

  1. Netty Server(接入层)
    • 门神:基于 NIO 模型,单机轻松抗住几万连接。
    • 协议解析:处理握手、拆包粘包(虽然 WebSocket 也就是帧处理)、心跳保活。
  2. Redis Pub/Sub(路由层)
    • 大喇叭:因为是集群部署,用户A连在服务器1,用户B连在服务器2。A发的弹幕,必须通过 Redis 广播给所有服务器,再推给B。这是分布式消息推送的核心!
  3. RocketMQ(削峰填谷)
    • 蓄水池:弹幕发太快,数据库写不过来怎么办?先扔进 MQ,让消费者慢慢存库。保证消息不丢,数据库不炸

💻 Round 3:核心代码秀(Show Me The Code)

3.1 Netty 启动类:性能调优的艺术

别光会 new ServerBootstrap(),参数调优才是关键!

@Component
public class NettyServer {
    
    @PostConstruct
    public void start() {
        ServerBootstrap bootstrap = new ServerBootstrap();
        bootstrap.group(bossGroup, workerGroup) // 主从Reactor多线程模型
            .channel(NioServerSocketChannel.class)
            // 握手队列长度,并发高时调大点,防止连接被拒绝
            .option(ChannelOption.SO_BACKLOG, 1024) 
            // 开启 TCP_NODELAY,禁用 Nagle 算法,有消息立刻发,拒绝延迟
            .childOption(ChannelOption.TCP_NODELAY, true)
            // 接收缓冲区,弹幕一般不长,32k够够的
            .childOption(ChannelOption.SO_RCVBUF, 32 * 1024)
            .childHandler(new ChannelInitializer<SocketChannel>() {
                @Override
                protected void initChannel(SocketChannel ch) {
                    ChannelPipeline p = ch.pipeline();
                    // 30秒没动静?直接踢下线,节省资源
                    p.addLast(new IdleStateHandler(30, 0, 0));
                    p.addLast(new HttpServerCodec());
                    p.addLast(new ChunkedWriteHandler());
                    p.addLast(new HttpObjectAggregator(65536));
                    // 处理 /ws 路径的 WebSocket 握手
                    p.addLast(new WebSocketServerProtocolHandler("/ws"));
                    // 咱们自己的业务逻辑
                    p.addLast(danmakuHandler); 
                }
            });
            // ... bind port ...
    }
}

3.2 弹幕处理器:敏感词退散!

弹幕里要是出现了“祖安语录”或者广告,平台是要担责的。所以我们引入了 DFA算法(Trie树) 做敏感词过滤,速度快到飞起。

@Component
@ChannelHandler.Sharable // ⚠️注意:Handler如果要复用,必须加这个注解
public class DanmakuHandler extends SimpleChannelInboundHandler<TextWebSocketFrame> {

    @Override
    protected void channelRead0(ChannelHandlerContext ctx, TextWebSocketFrame frame) {
        // 1. 反序列化
        DanmakuMsg msg = JSON.parseObject(frame.text(), DanmakuMsg.class);
        
        // 2. 敏感词过滤(DFA算法,效率比正则高N倍)
        // 如果是 "垃圾平台",会被替换成 "**平台"
        String cleanContent = sensitiveWordFilter.filter(msg.getContent());
        msg.setContent(cleanContent);

        // 3. 广播给所有人(先发到 Redis)
        redisTemplate.convertAndSend("danmaku_topic", msg);
        
        // 4. 异步持久化(扔给 MQ,别阻塞主线程)
        rocketMQTemplate.asyncSend("db_save_topic", msg);
    }
}

3.3 Redis 广播监听:跨服聊天全靠它

@Component
public class RedisMessageListener implements MessageListener {
    @Override
    public void onMessage(Message message, byte[] pattern) {
        // 收到 Redis 广播的消息
        DanmakuMsg msg = deserialize(message.getBody());
        
        // 拿到当前服务器所有活跃的 WebSocket 连接
        ChannelGroup allChannels = NettyConfig.getChannelGroup();
        
        // 群发!让每个连接的客户端都看到这条弹幕
        allChannels.writeAndFlush(new TextWebSocketFrame(JSON.toJSONString(msg)));
    }
}

💣 Round 4:踩坑实录(血泪教训)

开发过程不是一帆风顺的,我们踩过的坑,希望能帮你避雷。

😱 坑位1:消息丢了?学生要“寄刀片”

  • 现象:早起测试时,偶尔有几条弹幕凭空消失,尤其是在网络波动时。
  • 排查:Netty 的 writeAndFlush 是异步的,如果连接正好断了,消息就发丢了。
  • 填坑
    虽然弹幕允许少量丢失(没人会数是不是少了一条),但在点赞/送礼这种涉及数据的场景,我们引入了 ACK 机制。客户端收到消息必须回一个“收到”,否则服务器会重试。当然,对于普通弹幕,我们选择了“尽力而为”,毕竟为了100%可靠性牺牲性能不划算。

😵 坑位2:Redis 成了瓶颈

  • 现象:晚高峰大量弹幕刷屏,Redis CPU 飙高。
  • 原因:所有消息都走 Redis Pub/Sub,序列化和反序列化开销太大。
  • 填坑
    1. 消息合并:不要一条一条发,攒个100ms或者攒够10条,打包成一个 List 发送到 Redis。
    2. 精简协议:JSON 太啰嗦了,字段名死长。我们考虑过 Protobuf,最后为了调试方便,选择了精简 JSON 字段名(content -> c, userId -> u)。

🧟 坑位3:僵尸连接占着茅坑不拉屎

  • 现象:服务器显示在线人数10万,但互动的只有1万。
  • 原因:很多学生切后台、关机了,但 TCP 连接没正常挥手,服务器还傻傻维护着。
  • 填坑
    Netty 的 IdleStateHandler 是神器。心跳检测! 每30秒没收到客户端的心跳包(Ping),直接 ctx.close(),释放句柄,释放内存!

📊 Round 5:战果展示(拿数据说话)

上线一个月,看看这华丽的成绩单:

指标 表现 备注
单机并发 50,000+ 只要内存够,还能更高
消息延迟 < 200ms 人眼几乎无感
CPU 占用 峰值 45% 稳如老狗
用户体验 ⭐⭐⭐⭐⭐ “终于能看到老师回我了!”

最重要的是,点赞数比以前翻了10倍!看着满屏的爱心气泡,运营小姐姐笑开了花,技术团队的奶茶也有了着落。🥤


🚀 Round 6:画大饼(未来规划)

现在的系统虽然能打,但我们的征途是星辰大海:

  1. 多房间隔离:现在是全员广播,未来要支持几千个直播间互不干扰。这就需要在 Redis Topic 上做文章,topic:room:1001
  2. 礼物特效:光有文字不行,得有“火箭”、“跑车”。这涉及到复杂的动画指令同步和更严格的事务一致性(扣钱不能错!)。
  3. AI 智能审核:现在的敏感词库太生硬。未来接入 GPT,识别“阴阳怪气”的弹幕,自动禁言,净化网络环境。

📝 总结:技术服务于体验

这一路走来,从 HTTP 到 WebSocket,从单机到分布式,我们折腾技术的目的只有一个:让屏幕那头的学生不再孤单

技术不是冷冰冰的代码,它是连接人与人的桥梁。当看到第一条弹幕实时弹出,当看到全屏的点赞,所有的加班和掉发都是值得的。

下期预告
📝 智学云开发日记(六):薅羊毛大作战!如何设计防刷的营销优惠券系统?
想知道怎么用 Redis Lua 脚本秒杀高并发?点个关注,不迷路!😎


🔥 觉得有帮助的朋友,欢迎点赞、收藏、关注三连!评论区交流你的 Netty 踩坑经历!

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐