春节群聊背后的分布式消息队列:从家族群到红包高并发实战
大年初一早上七点我是被手机震醒的。家族群里二姨发了一段六十秒的语音表哥连着上传九张年夜饭照片小叔丢出一个红包紧接着满屏谢谢老板的表情包。消息像瀑布一样往下刷手机震了十几分钟都没停。我盯着聊天界面突然想到一个问题微信服务器凭什么能同时扛住几十个家族群在春节期间的暴力刷屏每个人发送的拜年消息、语音、视频、红包是怎么一条不漏地送到每个群成员手里的答案藏在一个很常见的后端架构概念里——分布式消息队列。你不需要把这事想得特别高大上文老师课堂-春节特别版这期就用家族群聊这个场景把分布式消息队列的核心机制、应用场景、常见坑位一次讲透。无论你是刚学后端、准备面试、还是在写CRUD时想知道中间件到底解决什么问题这一篇都能给你一套能直接拿去用的理解框架。看到最后我还会带你把一个简化版的消息队列Demo跑起来亲眼看一看生产消费、重复消费、消息堆积这些春节事故是怎么发生的。1. 家族群聊本身就是一个消息队列1.1 谁在发、谁在收、谁在中间存储把家族群拆开看它和消息队列的模型几乎一模一样。每个家族群成员既会发消息也会收消息所以每个人既是生产者Producer又是消费者Consumer。消息不会直接从一个人的手机飞到另一个人的手机而是先上传到微信的服务器再由服务器把这条消息推给群里的所有成员。那个服务器就是消息队列里的Broker代理/中间人。你可以这样对应家族群 消息队列里的Topic主题每个群聊都有一个独立的消息通道微信服务器Broker负责接收、存储、转发消息群成员消费者从自己的聊天窗口里接收消息私聊点对点队列一条消息只能被一个接收者消费群聊发布/订阅模式一条消息会被广播给所有订阅了该Topic的消费者。这里面有一个很多初学者容易忽略的点消息不是转瞬即逝的。你在群里发一条消息即使小叔当时手机没网等他连上WiFi打开微信消息还是会出现在他手机上。这就是消息队列的核心动作——持久化存储Broker把消息先存下来消费者什么时候来取、能不能取到由消费机制决定而不是发送方直接推到接收方手机里。1.2 你退群又进群为什么看不到历史消息家族群还有一个很生活化的现象刚进群的新人看不到历史消息退出再进群也看不到之前的聊天记录。放到消息队列里这其实就是**消费位点Offset**的管理问题。群成员相当于一个消费者微信服务器会给每个群成员维护一个已读位置。你看到哪一条了服务器就从这个位置之后继续给你推。新人刚进群时服务端没有为他建立历史消费位点所以默认从当前时间点开始推送。这和消息队列里消费者组从最新位置开始消费是一个逻辑。更妙的是未读消息数量。你几天没看群打开后看到99条未读消息这对应消息队列里非常关键的指标——消费滞后Lag。生产者已经写入了几百条消息而你的消费进度停在大年初一滞后量就是99条。运维监控消息队列时盯得最多的就是这个Lag一旦消费者的处理速度跟不上生产者Lag就会持续上涨最终造成消息堆积。再往深挖一点消息队列里有个能力叫回溯消费消费者可以指定从某个时间点或某条消息ID开始重新读取消息。家族群里下拉加载历史消息本质就是一次按时间回溯的消费。Kafka里可以按offset重放RocketMQ可以按时间回退Redis Stream里甚至可以按消息ID精确定位。原理一样都是把消息当作可回放的日志而不是消费完就销毁的一次性数据。所以你发现没有——你天天在用的家族群其实就是一张分布式系统的活教材。2. 用除夕抢红包说清楚消息队列为什么存在只看模型还不够很多人想不明白的是消息队列到底解决了什么问题为什么不能直接让A服务调B服务的接口这个问题用除夕夜的红包场景来解释再合适不过。2.1 削峰填谷全家同时抢红包服务器为什么没炸除夕零点家族群里所有人同时在发拜年消息、同时抢红包。假设这个群瞬间有200个人同时在操作社区团购群、公司群、同学群全加起来同一秒可能有几十万次请求打到服务器。如果不用消息队列架构往往是这样的用户点击抢红包——请求直接到后端服务——后端马上查数据库——扣减红包余额——写抢到记录——返回结果。这就是一次同步调用。如果数据库同时收到几十万条请求连接池瞬间耗尽数据库CPU飙升最后的结果就是红包页面转圈圈谁都没抢到。这时候消息队列出场了。它会先把我要抢这个红包这个请求立刻接收并存储然后返回给用户一个抢到了排队中的响应。真正的扣减逻辑放到后台由消费者按照自己能够承受的速度一条条处理。秒杀期间数据库每秒只需要处理几百个请求压力完全可控。这个动作就叫削峰填谷——把瞬时高峰流量接到队列里削平峰值再用稳定的速度慢慢消化。你可以把它理解成食堂打饭一个窗口突然涌进来200个人肯定挤爆但如果在门口设一个排队区域让队伍有序往里放人窗口的接待压力就被控制住了。消息队列就是那个排队区它不一定让单个请求处理得更快但能让整个系统在大流量下不崩溃。2.2 异步解耦拍祝福视频的人不用等所有人看完再来看另一个场景。你在群里发了一条拜年视频系统要做的事情不止把视频发出去要保存原视频、生成不同清晰度的转码版本、更新群聊动态、给所有群成员推送通知、记录播放数据……如果这些动作全部同步执行用户发一条视频可能要等好几秒钟期间任何一步失败比如转码服务超时整条链路都会失败。消息队列提供的是异步解耦发视频这个主流程只需要把我发了一条视频这个消息丢进队列立刻返回成功。视频转码服务、通知推送服务、数据统计服务各自订阅同一个Topic谁有空谁处理互不干扰。就算通知服务挂了视频已经发出去了等到通知服务恢复后再重新消费消息补通知就行。家族群里的群接龙也是这个套路。发起人只负责发布接龙开始这条消息并不需要站在原地等群里每一个人回复大家的接龙内容陆续回来后台再把名单汇总好。上游不用等下游下游之间也不互相依赖这就是解耦的直接好处。分布式架构里服务的拆分越来越细消息队列几乎成了解决服务间异步协作的默认选项。拍视频这步如果要存文件一般还会配上对象存储比如MinIO之类的分布式存储流程是先把文件传到存储再把文件已传好请转码的消息丢进队列转码服务收到消息后去存储里取文件。2.3 可靠投递与收到请回复群里最让人血压升高的消息是什么是收到请回复。长辈在群里发了一句今晚八点视频通话收到请回复本质是一条需要确认的消息。微信的已读回执就是ACK如果你一直不回长辈可能会私聊再问一次这就是重试如果私聊也没回长辈可能直接打电话这就是消息队列里的死信处理。把这套生活逻辑映射到技术上消息队列提供了几种投递语义投递语义含义生活类比At most once至多一次消息可能丢失但绝不重复消息发出去就行不确认收到没收到At least once至少一次消息保证送达但可能重复收到请回复没回复就重发可能发两次Exactly once精确一次每条消息只被消费一次精确到人、精确到次数的事后确认真实生产系统里最常用的是At least once。它实现简单性能好代价是消费者必须考虑重复消费。家族群里长辈发了两遍收到请回复你要回答两遍收到这在生活里是礼貌但如果是红包扣款消息被消费了两次就会把同一个红包发两次钱这个后果就很严重了。所以消费端要做幂等处理用红包ID去重同一笔红包只允许领取一次。那些重试很多次都没人消费成功的消息最终会被丢进死信队列DLQ。在生活里这就是群里了所有人也没人理最后只能单独打电话通知。死信队列存在的意义是让始终消费不掉的消息不要无限占用队列资源而是进一个专门的地方等人工排查或定时任务捡走处理。3. 手把手搭一个家族群聊版消息队列Demo概念听得再多不如动手跑一遍。接下来我用Redis Stream做一个简化版红包消息队列模拟长辈发红包、三个消费者进程抢红包的场景。如果你机器上没有Redis先别急跟着下面两步走3.1 选型为什么用Redis Stream而不是Kafka你可能会问学消息队列不是应该直接上Kafka吗我的建议是先用Redis Stream入门理由很直接Kafka部署要处理ZooKeeper或KRaft本地环境调试成本高Redis一条命令就能起Redis Stream天然包含消息ID、持久化、消费者组、ACK机制这些消息队列核心概念用Redis练明白了套路再去看Kafka、RocketMQ你会发现除了分区模型更复杂、性能更强本质逻辑是相通的。对比项KafkaRedis Stream部署成本高依赖组件多极低单机Redis即可数据量规模十亿级消息毫无压力适合中小规模、临时队列消费模型分区消费者组顺序性强消费者组适合轻量任务学习曲线陡峭平缓适用场景大流量日志、事件流、核心业务MQ轻量异步任务、限流缓冲、教学演示本地起一个Redis 7.x用Docker最省事docker run -d --name redis-mq -p 6379:6379 redis:7然后装Python的redis客户端pip install redis3.2 生产者长辈发红包的代码长什么样生产者要做的只有一件事向Stream里写入一条消息。Redis Stream的结构很直观一个Stream就是一套按时间排好的消息日志每条消息有一个自动递增的ID正文是多个字段组成的键值对。import redis r redis.Redis(hostlocalhost, port6379) STREAM family:redpacket # 大舅在群里发了一个100元的红包 msg_id r.xadd(STREAM, { sender: 大舅, amount: 100, redpacket_id: rp_20250210_001, target_group: 文氏家族群, }, maxlen1000) print(f红包消息已写入ID {msg_id.decode()})这里有几个细节值得说明redpacket_id是整个消息的业务唯一ID千万不能省。后面做幂等、防重复消费全靠它。maxlen1000是给Stream设置了一个长度上限超过1000条自动清理旧消息。这对应了线上消息队列的消息保留策略避免Broker无限堆积数据撑爆内存。msg_id是Redis自动生成的时间戳序号组合类似Kafka里的offset作用。在执行这段代码前先启动一个Redis的监控命令方便观察redis-cli -n 0 monitor你会看到XADD命令实时出现在监控窗口里说明消息真的进了Broker。3.3 消费者抢红包逻辑与ACK机制消费者这边要用的概念比生产者多一些。先创建消费者组再用组内成员去读取消息。import redis import time r redis.Redis(hostlocalhost, port6379) STREAM family:redpacket GROUP redpacket-grabber # 创建消费者组只创建一次 try: r.xgroup_create(STREAM, GROUP, id0, mkstreamTrue) print(消费者组创建成功) except redis.ResponseError as e: print(f消费者组已存在或创建失败: {e}) def consumer(name, count10): # 表示只读取从未投递给当前消费者的新消息 results r.xreadgroup(GROUP, name, {STREAM: }, countcount, block2000) if not results: print(f[{name}] 没有新消息) return for stream_name, messages in results: for msg_id, data in messages: # 模拟抢红包打印谁抢到了 print(f[{name}] 抢到红包: {msg_id.decode()} {data}) # 处理完成后确认ACK r.xack(STREAM, GROUP, msg_id) # 模拟处理耗时让消费者看起来更真实 time.sleep(0.1) # 第一次模拟消费者A去消费 consumer(uncle-li, count20)这段代码的关键在于三个参数GROUP消费者组名。同一个组里的消费者是竞争关系一条消息只会被组里其中一个成员拿到而且 Redis 会把不同消息尽量平均分配。特殊的读取标志表示只读从未被投递给当前消费者的新消息。如果要回溯历史消息可以指定一个具体的消息ID这就类似翻聊天记录。xack表示消费成功后向Broker确认。如果处理完没确认就宕机重启后这条消息会被重新投递给某个消费者——这正好解释了后面要讲的重复消费问题。3.4 多进程模拟三个家庭成员抢红包为了模拟分布式效果我们用三个进程同时让三个家庭成员抢红包。不想写多进程代码的话打开三个终端窗口、分别执行上面消费者函数也行逻辑等价。from multiprocessing import Process p1 Process(targetconsumer, args(爸爸手机,)) p2 Process(targetconsumer, args(妈妈手机,)) p3 Process(targetconsumer, args(小叔手机,)) for p in [p1, p2, p3]: p.start() for p in [p1, p2, p3]: p.join()运行生产者往队列里放10个红包再启动三个消费者进程你会看到类似这样的输出[爸爸手机] 抢到红包: 1739200000000-0 {sender: 大舅, amount: 100, ...} [妈妈手机] 抢到红包: 1739200000001-0 {sender: 二姨, amount: 88, ...} [小叔手机] 抢到红包: 1739200000002-0 {sender: 表哥, amount: 66, ...}三条进程一起跑红包被瓜分得很快。这就是消费者组的核心效果同一批消息组内多个消费者分工协作互不重复。这个互不重复由Broker保证了不需要你自己加锁协调——多节点分布式部署时这是巨大的便利。跑完之后可以用XPENDING看正在处理中的消息redis-cli XPENDING family:redpacket redpacket-grabber如果所有消息都确认了这个命令会显示0条待处理如果故意把某个消费者进程kill掉不让它执行XACK你就能在PEL待处理消息列表里看到它占着几条消息没还回去。这就为下一节的踩坑埋下了伏笔。4. 春节场景里最容易踩的四个坑光会跑Demo还不够下面这四个问题是任何消息队列在生产环境中都会遇到的春节事故现场我把排查思路和解决方案也一并写出来。4.1 重复消费为什么舅舅收到了两次红包通知现象除夕晚上你发了一个拼手气红包舅舅点开后提示手气最佳结果过了两分钟又收到一条同样的通知。根因消费者在抢红包流程里完成了扣减操作但在还没执行XACK确认时进程刚好崩溃了。Broker没有收到确认认为这条消息还没处理完等该消费者恢复后就把消息重新投递给它或者其他消费者。嗳——同一个红包又被处理了一遍。排查链路先确认是否存在重复查业务日志看同一个redpacket_id是否有两次成功领取记录再到Redis执行XPENDING family:redpacket redpacket-grabber如果PEL里有大量待确认消息说明确实存在消费了但没确认的情况检查消费者代码确认XACK是否放在了try...finally里是否在处理完成后一定能执行。解决方案消费端必须做幂等。最简单的方式是给红包领取表加redpacket_id唯一索引第二次插入直接报错或者用Redis的SET NX实现分布式锁用redpacket_id user_id作为锁的key抢到锁才允许扣款。线上项目一般会用Redisson这类封装好的分布式锁组件配合看门狗机制防止锁过期导致并发问题。注意这个幂等方案不能在单机内存里做。红包服务如果是多节点部署的每个节点的HashMap都是自己一份A节点处理过的订单B节点根本不知道照样重复扣款。共享的Redis锁或者数据库唯一约束才能保证多个服务实例之间的一致性。4.2 消息乱序家族群里的消息顺序为什么会错乱现象表哥先发了一句我到了后发了一张停车场的照片但群里显示的顺序却是照片在前、文字在后。根因多种可能同时存在。可能是表哥手机4G网络抖动第二条消息先到了服务器也可能是消息进了同一个Topic但被两个消费者并发处理处理速度快的那条先被展示出来。消息队列本身往往不保证全局顺序只保证分区内有序。Redis Stream的写入顺序是有序的但从写入有序到展示有序中间还隔着一层消费处理速度。排查链路确认消息的生产时间是否本身乱序服务器日志看时间戳确认同一业务对象比如同一个人的连续操作是否被分配到同一个消费者处理检查消费者的处理逻辑是否存在多线程并发调用同一个外部接口。解决方案你不需要保证所有消息全局有序只要保证同一个用户的操作有序就够了。做法是给消息指定分区键同一个user_id的消息每次都分发到同一个消费者只用一个线程去处理天然有序。在Kafka里指定同一个key的消息进同一个partition。在家族群接龙场景里非常实用的一个做法是数据库表里加一个seq序号字段消费端把消息接收下来先存库展示端按seq排序就算处理顺序乱掉最终库里排序结果也是对的。生活里接龙报名如果多人同时编辑表格顺序就会乱但最后按序号整理一遍名单就没问题——一个道理。4.3 消息堆积大年初二还有人收到除夕的拜年消息现象大年初一凌晨几百个群同时刷屏消息队列里的消息数量肉眼可见地往上涨。到了初二还有人在发出新年快乐的延迟通知。根因生产者一瞬间生产了海量消息消费者的处理能力不够Lag持续上涨。常见原因有两个一是消费者实例太少二是消费者的下游依赖太慢比如每条消息要同步调用一个耗时2秒的AI语音转文字服务。排查链路先看滞后量redis-cli XLEN family:redpacket和XINFO STREAM family:redpacket看Stream当前总消息数看每个消费者的处理耗时明细定位是不是下游RPC超时确认消息消费速率是否远小于生产速率。解决方案最简单粗暴的办法横向扩容消费者实例从3个消费者加到10个处理能力瞬间翻几倍优化消费链路把同步调用改成批量处理比如攒够50条消息一次性批量写库吞吐量能提升一个量级消息里如果带视频文件把视频转码后的耗时操作继续往下游丢不要在红包消费者里同步等待设置合理的过期清理策略Redis Stream用XTRIM限制长度Kafka按配置的消息保留时间清理。对于非核心消息比如谁又发了个表情包的统计类数据堆积太多时可以直接丢弃保障核心链路红包、转账类的消息单独用一个高优先级Topic闲杂通知走另一个低优先级队列。宁可表情包统计丢几条也不能让红包消息排队排到初二。4.4 分布式事务红包扣款和入账怎么保持一致现象你发了100元红包一个人抢走50元另一个人抢走30元但记账服务显示红包总额还剩30元而实际上已经被抢完了。根因发红包、抢红包、流水记录、通知统计通常由不同的服务负责每个服务有自己的数据库。跨多个数据库的更新本地事务已经管不住了这就进入了分布式事务的范畴。如果直接同步调用来保证一致性先调账户服务扣减、再调记账服务写流水、再调通知服务发消息任何一步失败都需要回滚所有服务复杂度极高性能也很差。现实中更常见的方案是基于消息队列的最终一致性。一套简化可落地的流程发红包服务在自己的数据库里先写入一条流水初始状态是待发送同一个本地事务里把红包已发送的消息写入消息队列记账服务、通知服务异步消费消息各自处理自己的业务消费成功后向Broker发送ACK标记完成如果某个服务消费失败消息进入重试队列重试重试多次仍然失败进死信队列由分布式定时任务比如XXL-Job定期扫描把待发送状态的红包重新发起补偿流程。这种方案要求你接受一个事实系统的状态不是瞬时一致的而是最终一致的。红包发出去的一瞬间账户可能还没扣款但几秒后一定能对上账。消息队列本身不解决原来那些同步调用需要严格回滚的问题但它给了你一个更宽松、更能容忍故障的时间窗口让复杂问题能在后台慢慢收敛。实际操作里事务消息比如RocketMQ的事务消息和本地消息表是两种主流实现Redis Stream本身不直接提供事务消息所以生产环境做资金类业务建议用成熟的RocketMQ或Kafka配合设计。官方文档里关于事务状态的半消息回查机制也建议认真读一读它解决的就是Broker不知道本地事务到底成没成功这个经典难题。跑完整个Demo又把这四个坑从头到尾捋过一遍之后我对消息队列的体会其实很朴素它就是一套把人与人之间、系统与系统之间的协作变得更有弹性的机制。想把这套东西用好关键不在背概念而在碰到线上问题时有清晰的处理顺序——先看Lag和消费速率再定位逻辑错误最后才考虑架构层面的改造。春节假期如果有空建议你也把这套Demo跑一遍亲手让消费者崩溃一次亲眼看看PEL里的pending消息比看十篇博客都管用。
上一篇/下一篇内容由系统自动关联
返回资讯列表 →