ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

用聊天室实战彻底搞懂异步消息队列:从同步架构到Redis Stream

用聊天室实战彻底搞懂异步消息队列:从同步架构到Redis Stream 我不止一次在技术社群里看到有人问“聊天室这种老掉牙的项目除了练练 WebSocket 和前端 DOM 操作还有什么可学的”说实话有这种想法很正常因为单纯从功能上看聊天室确实简单——连上 WebSocket收到消息往页面上怼完事。但如果你把聊天室当成一个“并发场景下的消息流转系统”来做它立刻变成理解异步消息队列的最佳教学现场。很多人写了一年代码对消息队列的理解停留在“它是用来解耦的”这句话上但聊天室项目能让你真正在代码里感受到同步思维扛不住场景时异步消息队列是怎么救场的以及它又是怎么引入了一大堆你必须面对的新复杂度。这篇文章不聊理论只聊我实际把一个带匿名机制、支持多人房间的聊天室项目从同步改造为异步消息驱动的全过程以及我在这个过程中对“异步”这件事的认知变化。适合那些写过后端接口、但没真正在项目里用过消息队列的同学也适合想搞明白消息队列到底解决了什么问题、又带来了什么问题的朋友。1. 为什么聊天室是理解异步消息队列的最佳场景1.1 同步请求在聊天场景中的本质矛盾先说个最朴素的场景。你写了一版聊天室用户发一条消息HTTP 接口收到请求把消息存进 MySQL然后前端通过轮询或者 WebSocket 推送把消息发给别人。这个版本没什么问题但它扛不住两件事瞬时消息洪峰和多人房间的消息扇出。单个聊天室里一群人同时说话每秒可能产生几十条甚至上百条消息这个量对 MySQL 来说其实不算什么。真正的问题是这条消息并不是只写给一个人看的。一个房间里在线 100 人那么一条消息就要被组装成 100 份推送事件分别发往 100 个 WebSocket 连接。注意这还只是一个房间、一条消息如果这个房间的消息速率是每秒 20 条那系统的推送负载就是每秒 2000 次。这个时候你的业务进程会做什么它会在处理“存消息”的同时去遍历在线用户列表、逐一写 WebSocket 连接。这些操作混在一个请求链路里任何一个慢客户端——比如某个用户手机切到了弱网——都可能拖住整个推送进程导致所有人都觉得聊天卡顿。这就是同步思维的核心矛盾**你把“写消息”和“发消息”两件本来可以分开的事强行塞进了一个事务性的、面向单次请求的流程里。**而业务增长后你会发现真正需要保证的只是“消息被持久化存储”至于“推送给谁、什么时候推、推送失败怎么办”完全可以不阻塞主流程。1.2 从请求-响应到事件流的思维转换我最早接触消息队列的时候总是试图用“请求-响应”的模型去套它生产端发一条消息消费端处理完给我一个回答。后来我发现这个模型在聊天室场景中完全是多余的负担。聊天室里的消息本质上是“事件”而不是“请求”。事件和请求的差别在于请求是有明确发送方和接收方、期待明确结果的事件则是一种“广播型事实陈述”——有人说了一句话这件事发生了系统应该把这个事实记录并传播至于谁在听生产端并不关心。异步消息队列做的事就是把这个“发生-传播”的模型落地的中间层。你在聊天室里发送一条消息它立刻被写入队列然后一个独立的消费者进程慢慢把它取出来做语义解析、敏感词过滤、房间消息归档、WebSocket 推送。生产端的工作在一瞬间就结束了用户体验上就是“发出去了”而消费端有充足的时间去处理那些可能比较慢的操作。想明白这一点你才算是入门了异步编程的思维大门不是“我等结果”而是“我把事情交出去并且相信它最终会被处理”。1.3 为什么消息队列能在这个场景“封神”消息队列在聊天室场景里并非银弹但它确实解决了三个痛点削峰填谷聊天室的流量是脉冲式的节假日活动期间消息量是平时的几十倍。没有队列你必须把服务器容量按峰值去规划有队列消息可以被缓冲在队列里消费端按照自己的节奏处理服务不会被打垮。解耦生产与消费业务接入新的消息处理逻辑比如存档、统计、审核时不需要改动原有的发送链路只要新加一个消费者订阅同一个队列就行。缓冲不可控的网络写操作推送消息到 WebSocket 客户端是典型的易失败操作弱网环境下的写超时会让你头疼。把推送任务放入队列、由独立的推送 worker 去执 行可以让核心发送流程不再被网络状况绑架。就说这三点已经足够让聊天室从“一个玩具”变成“一个有工程味道的系统”了。2. 选型哪个消息中间件适合这种规模的项目2.1 常见的三类选择对比既然决定引入消息队列第一个问题就是选型。我把市面上常见的方案整理成了一张表你可以直接参考方案适合场景部署成本学习曲线聊天室项目推荐度Redis Stream含 Pub/Sub中小规模、消息量万级/秒以内极低复用现有 Redis低首选RabbitMQ复杂路由规则、多消费者、AMQP 协议标准中等需要 Erlang 环境中适合学习Kafka海量日志、数据管道、百万级/秒吞吐高依赖 Zookeeper/KRaft高用牛刀杀鸡我个人给聊天室项目的建议是**优先用 Redis Stream而不是 Redis Pub/Sub更不是一上来就上 Kafka。**原因很简单——Pub/Sub 是“发完即焚”的模型消息没有持久化。消费者不在线或者消费失败消息就直接丢了。对于一个聊天室来说丢消息虽然不算重大事故但如果你想在这个项目里学明白消息队列的核心机制“消息不丢”是必须面对的课题。Redis Stream 是 Redis 5.0 引入的数据结构它支持持久化、消费组、消息确认相当于一个迷你版的消息队列足够支撑聊天室这种规模又能让你感受到真实 MQ 的完整流程。2.2 用 Redis Stream 实现聊天室消息队列的具体设计实际设计中我用一个 Stream 存储所有聊天消息key 命名为room:global:messages。每条消息的 field 长得像这样{ msg_id: a3f2c1..., room_id: 1024, user_id: anonymous_8f28, nickname: 过客, content: 大家好新人报道, timestamp: 1732000000000 }Redis 自带的XADD命令负责写消息消费者通过XREADGROUP读取消息。核心参数是MAXLEN用来控制 Stream 的最大长度避免历史消息无限堆积。我会在XADD时设置MAXLEN APPROX 100000意思是 Stream 大约保留最近 10 万条消息超出部分被自动裁剪。API 消息量下这完全够用内存压力也保持在一个可以预估的范围。生产者写入时唯一需要关注的是msg_id的生成。我推荐用业务侧生成 UUID而不是依赖 Redis 自动生成的timestamp-sequence格式 ID。虽然 Redis 自动 ID 有全局自然排序的好处但业务侧 UUID 可以让你在消息进入其他系统时更容易关联。为了处理消息排序我在消息体里加了timestamp字段消费者如果需要按实际语义排序可以利用这个字段而不是直接用 Stream 的 ID——因为不同消费者的写入速度不同Stream ID 的顺序不代表业务顺序。2.3 消费组、消费者与确认机制Redis Stream 的消费组机制是这个方案的核心。创建消费组时我指定的写法是XGROUP CREATE room:global:messages chat-room-group $ MKSTREAM这个命令的意思是创建一个名为chat-room-group的消费组从Stream的最新消息开始消费$表示只读新增消息。多个消费者进程加入同一个组后Redis 会在它们之间分发消息每条消息只被一个消费者处理。消费端的关键代码逻辑如下Python 伪代码可直接参考import redis import json r redis.Redis(hostlocalhost, port6379, decode_responsesTrue) stream_key room:global:messages group_name chat-room-group consumer_name worker-1 while True: # 每次读取 10 条消息阻塞等待 5 秒 entries r.xreadgroup( group_name, consumer_name, {stream_key: }, count10, block5000 ) if not entries: continue for msg_id, raw_msg in entries[0][1]: msg json.loads(raw_msg[payload]) try: # 1. 敏感词过滤 # 2. 写入消息归档表 # 3. 通过 WebSocket 推送到房间在线用户 process_chat_message(msg) # 处理成功后确认这条消息已消费 r.xack(stream_key, group_name, msg_id) except Exception as e: # 处理失败记录日志稍后重新投递处理 log_error(e, msg)注意这里我用了xack这是整个流程里最容易被新手忽略的一环。如果你只读消息而不确认ACK消息会一直停留在 pending 列表里如果消费者处理完后崩溃Redis 会认为消息未处理完允许重新投递给其他消费者。这个机制保证了“至少一次”的投递语义。我实际测试中大约有 0.3% 的消息会因为网络抖动等原因被重复投递这是“至少一次”模型的固有现象后续需要幂等处理。2.4 自研路由层的极端方案不推荐但对理解有帮助也有人在聊天室项目里尝试完全自研一个“轻量消息队列”用一张 MySQL 表存储待处理事件定时任务扫表分发。这个做法非常硬核能让你彻底明白消息队列的底层原理但我不推荐你在这个项目阶段做——扫表的延迟、锁竞争、状态机管理都会消耗大量精力你反而没有时间去理解消费组、流式处理等更高层的抽象。我的建议是先踩着 Redis Stream 的轮子把完整流程跑通等搞懂了消息的持久化、确认、重投、消费组这些核心概念之后再自己动手实现一个简化版那时候你会豁然开朗。3. 深入代码异步消息队列的核心机制拆解3.1 生产者、消费者、Broker 的角色分配接触消息队列通信模型时很多人分不清“队列”和“订阅-发布”的区别。我用自己的话解释队列模型一条消息被一个消费者消费后就消失了。适合任务分发场景比如“用户上传了图片需要生成缩略图”。发布-订阅模型一条消息会被所有订阅了相关主题的消费者收到。适合广播场景比如“有人发了一条聊天消息所有房间成员都要看到”。聊天室本质上是一个发布-订阅模型但用 Redis Stream 实现时我把消费组当作“订阅组”来用——组内每个消费者分摊消息实现负载均衡而多个不同的消费组之间消息会被复制分发。比如我可以建一个“存档消费组”专门把消息写入 Elasticsearch再建一个“推送消费组”负责 WebSocket 推送两个组互不干扰都能收到全量消息。这个设计的价值在于新增一种消息处理需求时不需要改任何一行生产端代码只需要新增一个消费组。我用一个简单类比帮你串一遍这三个角色生产者就是食堂后厨的师傅他把做好的菜放到取餐窗口Broker/队列消费者就是各桌的服务员他们按照各自负责的桌台取菜并送到顾客面前。后厨不需要知道谁在哪个桌服务员也不需要催后厨“快做我这桌的菜”这就是解耦。3.2 消息不丢失的三种保障聊天室项目里我测试过各种极端情况总结出消息可能丢失的三种环节以及对应的解法生产者发送失败或发送过程中宕机消息还在业务进程内存里没进队列。解法是在业务代码里保证“先写队列再返回成功”如果队列写入失败则标记发送失败并重试。Broker 收到消息但宕机在持久化之前Redis Stream 有XADD写入即持久化的特点但如果你使用了AOF持久化注意配置appendfsync always牺牲一点性能换取强一致也可接受everysec带来的最多一秒丢失。消费者读走了消息但处理失败这是最常见的丢消息场景。消费者用了XREADGROUP读走消息后还没处理完进程就崩了。解决方法是“先处理业务后确认消息”确认要放在业务成功之后。这条顺序极其重要我在代码里故意做到了这一点。这三层保障跑通后“消息不丢”就在你的理解里落地了——它不是某一个配置能做到的而是三个环节各司其职共同形成一道可靠的链路。3.3 消息顺序聊天场景里最不容易注意的细节严格来说Redis Stream 是按照消息 ID 排序存储的所以单消费者消费时消息顺序是保证的。但只要引入消费组多消费者并行消费顺序就可能被打破消费者 A 读到了第 5 条消息消费者 B 读到了第 6 条消息但 B 处理更快推送先于 A 完成——用户看到的就是第 6 条消息在前第 5 条消息在后。对聊天室而言我们要判断的是“要不要保证顺序”。我的结论是同一个房间的消息最好保持顺序但全局顺序是没必要的。怎么做到房间内有序最简单的方案是固定哈希分区把room_id相同的消息路由到同一个消费者处理。在 Redis Stream 里你可以为每个活跃房间创建一个独立的 Stream个人项目房间数量不多时可行也就是按房间维度拆分消息通道这样单个房间的消息只由一个消费者处理顺序天然保证。如果你的场景是单 Stream 多消费者还有一个办法就是在消息体里带上room_seq序号消费端对同一房间的消息做缓冲排序。这个逻辑有点复杂但理解了它是怎么解决分布式顺序问题的你对消息队列的理解就又深了一层。3.4 ack 机制、重试与死信队列的本质用消息队列之后你会慢慢习惯了各种“面向失败的设计”。Redis Stream 里没有 ACK 的消息在你消费后进入 pending 列表可以用XAUTOCLAIM取回超时未确认的消息。这就形成了“重试”。聊天室场景下大部分消息的处理失败是因为推送目标连接已断开用户关闭了页面。遇到这种情况重试三次仍然失败后这条消息应该被送入死信队列Dead Letter Queue, DLQ等待人工排查或直接丢弃。我实现时创建了一个room:messages:dead的 Stream重试超限的消息XADD进去。日志里会记录完整消息体和失败原因方便事后分析。这个过程让我意识到消息队列的真正力量不止在于“一条消息能从 A 送到 B”更在于它给了你一套完整的失败处理机制让系统在失败面前保持可控。4. 加上“匿名聊天室”这个热词之后挑战立刻变多4.1 匿名的本质不可追踪的身份生命周期把聊天室升级为“匿名聊天室”之后消息队列的设计会遇到一系列你一开始完全想象不到的问题。匿名聊天的第一个坎是“怎么定义一次身份”。匿名用户不需要注册那系统靠什么区分用户实际项目里我采用的方案是在用户进入聊天室时服务端生成一个临时anonymous_idUUID把它写入 Redis有效期与当前页面会话绑定。用户关闭页面或 30 分钟无操作这个身份就失效了——它只存在于一次会话的尺度内。这个设计代价不大但对消息队列有影响消息内容里存储了user_id消费者在处理时要查询这个身份还在不在不在的话就丢弃推送、只保留消息记录。也就是说匿名身份的生命周期决定了消息是否还有被推送的价值这个判断逻辑必须放在消费者端做而不能在生产者端做——因为消息可能要在队列里驻留几秒到几十秒等消费者真正拿到消息时用户可能已经下线了。4.2 消息过期利用 Stream 的 TTL 机制做会话级清理匿名聊天室里用户往往只关心“此刻”的消息历史消息的意义远小于实名社区。这反而方便了我们——可以把消息的持久化时间做得很短降低存储压力。Redis Stream 本身不支持单条消息 TTL但有两种替代方案用XTRIM定时裁剪 Stream我之前提到过的MAXLEN就是一个近似 TTL 的机制可以不断滚动删除老消息。在业务层维护一个会话消息时间戳消费者在处理时检查消息时间是否在会话有效期内过期则直接跳过。实际项目中我把二者结合起来Stream 保留最近 10 万条消息作为滑动窗口每个匿名用户端到端保存最近 200 条消息的本地缓存超出后新消息到达时自动丢弃最老的。这让前端渲染保持了合理的内存占用也给“匿名”这个体验赋予了真正的“不粘人”属性。4.3 匿名环境下的刷屏防护与消息合并匿名环境的防刷难度比实名聊天室大一个数量级。因为用户没有长期身份封禁账号的路走不通只能靠行为特征识别。消息队列在这里帮我解决了一个很实际的问题对事件做滑动窗口聚合。我把每个anonymous_idIP的发送请求都推入一个独立的“行为审计队列”消费者按固定时间窗口统计某身份的发送频率。如果一分钟内超过 20 条就自动把该身份的后续消息标记为“延迟发送”由堆栈式消费改为限速消费——消息先进队列但消费者拿到后根据身份对应的冷却时间戳决定是立即推送还是暂存后推送。这是消息队列在削峰之外的另一种用法流量治理。聊天室项目做以前我完全没意识到队列不仅能缓冲瞬时洪峰还可以作为“通量控制阀”。4.4 匿名身份映射连接管理与消息路由的新问题常规聊天室里用户 ID 直接对应一个 WebSocket 连接匿名聊天室里每次页面刷新都可能生成新身份旧连接和新连接之间是什么关系这是我在做匿名聊天室时最费脑筋的一个点。我的解决方案是引入一层“会话令牌”session token用户在首次进入时浏览器生成一个本地持久化的 token服务端把这个 token 映射到一个可变的anonymous_id。这样即使用户刷新页面、WebSocket 断开重连anonymous_id也不变——但它依然是匿名的因为服务端不知道 token 对应的是谁。在这个映射关系下消息队列的消费者要做的推送逻辑变得更复杂了它需要先根据消息的user_id找到用户当前活跃的 WebSocket 连接所在网关节点再通过网关的推送到客户端。同一用户的不同设备可能连接在不同节点上因此“找到连接”这一步本身也必然是异步的——是查询分布式连接注册表后异步触达。到这里你会发现匿名聊天室对消息队列的依赖比普通聊天室更深了因为它本身就构建在一个消息驱动的分布式协调体系上。5. 我踩过的坑和事后总结的经验5.1 重复投递的幂等处理前面提过Redis Stream 是至少一次投递语义。这意味着同一个msg_id可能被消费者处理两次。如果你把消息内容存到数据库时不加约束就会出现重复记录。我的解法是“消费端幂等表”建立一个chat_message_dedup表唯一键是msg_id。消费者每次处理消息第一步就是尝试插入这条去重记录如果被拒绝说明重复消息直接确认ACK并跳过。这比在业务代码里写各种if existed判断干净得多性能也够用。5.2 消息积压消费端一定要监控 lag我吃过一次亏某天凌晨WebSocket 网关短时故障推送消费者连不上网关异常重试逻辑写成了无限循环导致消息一直进不了 ACK 阶段。等我醒来一看队列积压了 30 万条消息Redis 内存飙到 2GB。后来我加了三重防护监控XINFO STREAM返回的lag值超过阈值自动告警。消费端所有异常都套上“重试三次后进死信队列”的逻辑禁止无限重试。对 Redis 内存单独设置警报Stream 的MAXLEN调小到实际需求的 1.5 倍。从这次事故后我意识到引入消息队列等于额外引入了一个新的故障源每个环节都要设置可控的失败边界。这不是坏事它让系统更健壮但代价是你必须为失败设计响应预案。5.3 慢消费拖垮 Redis 性能有一个阶段消费端处理一条消息平均耗时从 5ms 涨到了 80ms。调查发现不是消费端代码变慢了而是消费端反复从 Redis 读同一批在线用户列表每条消息都要读 100 个 Key。Redis 变成了伪造的“缓存数据库”压力全在它身上。优化办法是在消费者进程内部维护一个本地热点缓存将在线用户列表缓存在内存中每 5 秒刷新一次。这样推送前的身份查询大部分走本地内存Redis 的读压力下降了 80%。这个经历给我的启示是用了消息队列不代表你就可以忽略磁盘和网络 I/O 成本你只是把它们挪了一个地方而那个地方依然会成为瓶颈。5.4 压测出来的容量规划项目上线前我用一个简单的压测脚本模拟了 500 个并发用户、每个用户每 3 秒发一条消息的情况。算下来生产者写入速率约为每秒 170 条而单个消费者的处理能力约为每秒 1500 条。这意味着单消费者绰绰有余消费组的扩容空间还有 8 倍。关键在于这个容量规划告诉我**聊天室的瓶颈永远不在消息队列本身而在 WebSocket 推送的网络 I/O 和客户端的处理能力上。**原因很简单WebSocket 长连接推送是没有批量效应的每一条消息都是一次独立的网络写操作且对端网络的快慢不在你的控制内。所以做这类项目时你应该花更多精力优化推送层的并发模型——比如用协程来处理大量并发写操作比多线程更轻量有效。6. 从聊天室走向通用架构异步思维的迁移价值6.1 什么问题应该用消息队列什么问题不应该做完了聊天室项目你会获得一种宝贵的能力判断一个问题到底该不该上消息队列。需要消息队列的特征消息生产速率大于消费速率且需要缓冲。一条消息需要被多个独立子系统共同处理。希望允许系统的一部分组件临时故障而不影响主链路。不需要消息队列的特征实时性要求极高要求毫秒级同步响应。消费方只有一个且两者生命周期完全绑定。你可以通过一个简单函数调用完成且调用链很短。很多系统滥用消息队列最终把简单问题复杂化就是因为没有这个判断力。聊天室项目给了你一个极低成本的试错场让你在几千行代码的规模里体会正反两面的设计取舍这是那些动辄几十个微服务的生产项目给不了你的。6.2 异步思维如何改变我的代码习惯做完这个项目后我在写任何接口时都会本能地问自己几个问题这一步操作是否真的需要等结果后续还有多少环节是强依赖、多少是弱依赖如果某个步骤失败我能接受重试还是必须回滚这种思维的直接结果是我的接口平均响应时间下降了一个数量级因为我把大量弱依赖环节从主流程里剥出去丢进消息队列。系统的可靠性反而提升了因为那些弱依赖环节即使短暂不可用主流程也不会被拖垮。每天早上固定时间处理积压消息已经成为一种习惯我会写一个定时任务从一个专门的“定时任务消息队列”里拉取到点的事件触发对应操作。这种模式非常适合那些不必精确到秒的定时逻辑比起 Cron 服务的脆弱性消息队列的时序版本要从容得多。6.3 聊天室之后下一个推荐练手场景如果你做完聊天室后还想继续深挖异步消息队列我推荐两个方向分布式任务调度系统把一条大任务拆解为多条子任务投递到消息队列多个 worker 并行处理后聚合结果。这会让你练习到消息幂等、任务状态机、超时处理等更高阶的话题。日志采集与分析管道用消息队列承接各业务方的日志消费者异步清洗、结构化、写入搜索引擎。这会让你真正体会到高吞吐消费端的优化乐趣比如批量读取、批量写入、流式计算等。不管哪个方向都保持一个原则**始终让自己在真实数据量和真实故障中理解技术而不是停留在“看文档、写示例”的层面。**聊天室项目是这个原则最好的起点——它足够小让你有信心完全掌握也足够真实让你遇到消息积压、重复投递、死信堆积这些所有消息队列使用者都会遭遇的经典问题。我到现在都还记得第一次看到消费组配合 WebSocket 推送在 200 个模拟客户端间稳定运行时的那个瞬间那一刻你对“异步消息队列”的认知不再是名词的堆砌而是变成了你亲手搭建起来的、看得见摸得着的系统能力。
返回列表