消息队列
消息队列
主要有以下工具:
- Kafka、Kafka中文文档:吞吐量高,延迟低,大数据、日志收集等场景
- RabbitMQ、RabbitMQ中文文档:简单易学,生态活跃,中小规模应用场景
- RocketMQ:阿里出品,支持事务消息,适合电商金融等可靠性要求高的场景
为什么要用消息队列?-> 典型场景:高并发冲击
比如电商双十一秒杀等活动,导致数据库超载,从而引入消息队列
经典处理方式(同步):
- 校验库存
- 扣减库存
- 创建订单...
如果用户过多,系统直接摆烂。。。
消息队列
作用:
- 异步处理:将用户的下单请求先写入消息队列,后端服务按需消费,避免直接冲击数据库。
- 流量削峰:通过队列暂存请求,平滑处理高峰期流量,保护数据库和核心服务。
- 解耦系统模块:订单服务、库存服务等可通过消息队列通信,降低系统耦合度,便于扩展和维护。
- 提升可用性与容错能力:即使某一环节暂时不可用,消息仍可暂存在队列中,防止数据丢失。
就像一个快递驿站,快递员放下包裹就去处理别的快递(对应生产者),你收到快递到达消息去取快递(消费者);如果包裹很多,也可以保存包裹,等待快递被取走。
异步:快递员放下包裹就去处理别的任务(不用等待所有操作完成,快速获得相应)
削峰:驿站可以保存包裹,等待你取走(后端根据自己的能力处理)
解耦:不用关心是谁送的(服务间不直接依赖,系统可灵活扩展)
消息队列 vs 经典处理
经典处理:用户发送请求后,完成所有相关的业务逻辑与同步操作,才返回信息
使用消息队列:用户发送请求后,校验请求,并完成相关操作,同时向消息队列投送请求,即可返回结果;接下来后台服务处理消息队列中的请求,可以根据后台自己的能力慢慢处理请求
抢购服务只管向消息队列发送消息,不用关心谁来管理库存、创建订单,如果有扩展需求,也可以增加模块,只需要去监听消息队列即可,不用修改原先的抢购服务模块
RabbitMQ(建议Docker拉取)
- 消息队列的基本概念与核心作用
- 使用 RabbitMQ
- 学习主流队列模型和应用场景(SDK开发包)
- 了解生产环境(消息可靠性,幂等性,高可用集群,死信队列,延迟队列,优先队列等问题)
- 深入底层原理
⚠️重点:多种工作模式
Simple 模式:一对一关系
一个生产者,一个消息队列,一个消费者
就像老板把任务派发给员工,队列是存储任务的容器,老板把任务放进去,员工从里面取出来完成。
❓️如果有很多个任务要处理,但是一个员工完不成,怎么办呢?-> 多找几个员工帮忙就行了,就像工作队列模式
工作队列模式:一对多关系
一个生产者,一个队列,多个消费者
老板发布了一堆任务,RabbitMQ依次把任务分发给员工,一个任务只能由一个员工完成,大大提高了效率
❓️但是这老板又要求每个员工写个工作总结,怎么办?-> 发布订阅模式
发布订阅模式:一对多关系
一个生产者,一个队列,多个消费者
每个员工有自己的工作队列,而不是共享一个队列,老板通过交换机来控制任务的分配,将任务分配给绑定的队列(比如广播交换机)
❓️那么如果老板想把某些任务发给特定员工呢?-> 路由模式
路由模式:一对多关系
一个生产者,一个队列,多个消费者
给每个员工的模式设计自己负责的路由键,老板发布任务时会指定路由键,由直接交换机根据路由键精确匹配,把任务发给对应的队列
❓️那么如果老板想把前端任务发给前端员工,后端任务发给后端员工呢,一个一个指定路由键太麻烦。 -> 主题模式
主题模式:一对多关系
一个生产者,一个队列,多个消费者
使用主题交换机,支持通配符进行模糊匹配,可以为每个员工队列指定通配符路由键,类似正则表达式,接受所有匹配的任务
还有 RPC(远程调用模式)等。。。
生产环境中的坑🕳️
RabbitMQ 服务器重启,队列里的消息会丢失吗?
会丢失。 -> 解决:保证消息不丢失,做好持久化(将消息从内存保存到硬盘)
持久化三件套:1. 队列持久化 2. 消息持久化 3. 交换机持久化
做好持久化,消息就不会丢了?
还要做好消息确认机制保证消息的可靠性
消息确认机制:
- 生产者确认:生产者发送消息后,要等待 RabbitMQ 的确认回复
- 消费者确认:消费者处理消息成功后,要通知队列(千万不要用默认的自动确认)
如果网络抖动或者消费者重启,可能会重复消费消息?
保证消息幂等性,让重复消费的消息等同于一次消费
幂等性设计是分布式系统的基本功
常见三种方案:
- 每条消息唯一 ID,消费前检查此 ID 的消息是否消费过
- 利用数据库的唯一索引
- 利用 Redis 分布式锁(不推荐,成本较高)
消息乱序怎么办?
如果要严格保证顺序,那就是用单队列 + 单消费者,但是性能很低
所以根据业务需求权衡一致性和性能,如果全都要,考虑 kafka 的分区机制,将同一个用户的消息路由到同一个分区,同一个分区内的消息严格有序,不同分区内的消息并发处理
消息处理失败怎么办?
使用死信处理机制:如果消息被消费者拒绝、消息过期、队列满了等情况会产生死信,可以配置死信交换机路由到死信队列,由专门的消费者监控死信队列,发现死信后,可以重试,告警,人工处理等
RabbitMQ 服务挂了怎么办?
使用主从集群 + 哨兵模式,也就是搭建集群 + Quorum 仲裁队列。
Quorum 仲裁队列:将数据自动同步到多个节点,保证数据一致性,主节点挂了,自动选取新的主节点
如果数据量太大,单个队列存不下,可以使用 sharding 插件创建分片队列,把消息分散到多个节点(RabbitMQ不擅长此场景)
RabbitMQ 的高级特性
- 延迟队列:死信队列 + TTL,比如自动取消超时订单
- 优先级队列:设置队列最大优先级,发送消息时指定优先级,比如付费用户优先处理
- Stream流式存储:消息消费后不删除,可保留历史、重复消费、回溯,比如实时数据分析,审计日志等
带着问题学习
- 消息路由机制
- Exchange 用什么算法匹配
- Binding 如何存储
- 队列存储结构
- 消息在队列还是磁盘
- 不同队列类型有什么区别
- AMQP协议
- 客户端和服务端怎么通信
- 数据格式是什么样的
- 持久化的实现
- 数据怎么写入磁盘
- 如何保证不丢失
