编程导航消息队列话题讨论

消息队列

24 参与
分享

快来分享你的内容吧~

点击登录,快来和大家讨论吧~
表情
图片
话题
打卡
综合
交流
文章
问答

消息队列从入门到跑路,保姆级教程!傻子可懂

你是小阿巴,刚刚为电商系统的双 11 大促开发了秒杀抢购功能。 0 点秒杀开始,每秒上万个用户同时点击抢购按钮,你的数据库瞬间被打垮! ![](https://pic.code-nav.cn/post_picture/1601072287388278786/au8fjnOdsDt6d5I1.webp) 你急得满头大汗,只能找到 “后端之狗” 鱼皮求助:阿巴阿巴…… 鱼皮看了看你的代码:用户每次抢购,你的代码都要同步处理校验库存、扣减库存并创建订单等多个操作。同时抢购的用户多了,系统就会压力山大。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/5KQWazMmidAf9VOy.webp) 你急了:那咋办啊?我总不能拒绝用户吧! 鱼皮嘿嘿一笑:可以用消息队列呀。 你一脸懵:消息队列?那是啥? ⭐️ 推荐观看视频讲解,更通俗易懂:https://bilibili.com/video/BV12qmyBQEwL ## 第一阶段:认识消息队列 鱼皮:消息队列(俗称 MQ)就像一个快递驿站。快递员作为 **生产者**,把包裹(消息)放到驿站(消息队列),收件人作为 **消费者**,自己到驿站去取。 这样一来,快递员放下包裹就能走,不用费时间等你签收(异步)。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/YalVnB3sQ25IeRUg.webp) 你也不用关心是谁送的、不用和快递员见面,有空去驿站取就行(解耦)。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/H4Tca4GIb3x1F0br.webp) 就算双 11 包裹很多,驿站也能暂存起来等消费者慢慢取(削峰)。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/rZCQYTzwqs7YuRht.webp) 这就是消息队列的三大作用:**异步、解耦和削峰。** 对于你的秒杀抢购系统,现在用户每次抢购都要同步执行各种操作(校验库存、扣减库存、创建订单),全部完成才能返回结果。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/PbNZRvskQJGRDerf.webp) 如果使用消息队列,用户点击抢购后,系统先做基本的校验并在缓存中预扣库存,确认可以抢购后,把抢购请求(消息)投入消息队列,就可以立刻返回结果了(页面)。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/KMVX6iWCJqkFijyk.webp) 接下来由后台服务从队列中取出消息并处理那些耗时的操作,比如创建订单、在数据库中真正扣减库存。 这就是 **异步**:用户不用等所有后端操作完成,就能快速得到响应。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/4BwduPtEJTM8kYEA.webp) 而且,抢购服务只管发消息,不用关心是谁来处理库存、谁来创建订单。假如以后要加短信通知、用户行为分析等功能,只需要让新服务也去监听消息队列就行,完全不用修改抢购服务的代码。 这就是 **解耦**:服务之间不直接依赖,让系统能灵活扩展。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/pI8PANEIKqvsjosq.webp) 如果同时抢购的用户过多,数据库处理不过来,也可以用消息队列做缓冲。所有的抢购请求先快速存入队列,后台服务可以根据自己的能力从队列中慢慢拉取消息并进行处理。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/nqhefTZasUrvin4R.webp) 这就是 **削峰**:削平流量洪峰,保护后端系统不被压垮。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/WXZK6GHVbhcLatJO.webp) 你恍然大悟:原来如此,消息队列也太强大了吧!俺要学,俺要学! 鱼皮:很有干劲嘛!主流的消息队列实现技术有: - RabbitMQ:简单易学、生态活跃,适合中小规模应用 - Kafka:吞吐量高、延迟极低,适合大数据、日志收集场景 - RocketMQ:阿里出品、支持事务消息,适合电商金融等对可靠性要求高的场景 你:阿巴阿巴,这么多都要学吗? 鱼皮:新手建议从 RabbitMQ 开始入门,学好一个再学其他的就很简单了。下面我就以 RabbitMQ 为例带你快速掌握消息队列必知必会的技术。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/3LbT2qXbEx6BI7i5.webp) 你:期待住了,这不得点赞狠狠支持?! ![](https://pic.code-nav.cn/post_picture/1601072287388278786/rQIxU9iqfHMzTXgt.webp) ## 第二阶段:RabbitMQ 快速上手 鱼皮:首先来安装 RabbitMQ,笨办法是先安装 Erlang 语言环境,再去 [官网下载](https://github.com/rabbitmq/rabbitmq-server/releases) MQ 的安装包并执行。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/cUeiIuuFSR4ugfgH.webp) 但是更建议直接用 Docker 容器技术,快速安装和运行,不用担心版本不兼容。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/IcImoSKKlUQgrKYK.webp) 安装完之后,访问 http://localhost:15672 就可以打开 RabbitMQ 内置的可视化管理界面(默认用户名和密码都是 guest),你可以在这里查看队列的运行情况、手动发送和接收消息。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/vWkzvEyQ3zSeKRWM.webp) 另外,RabbitMQ 也提供了命令行工具 rabbitmqctl,主要用于运维管理,比如创建用户、查看状态等。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/1IpuxwGDTdofsWs8.webp) 你:那我怎么用代码操作 RabbitMQ 呢? 鱼皮:RabbitMQ 提供了各种编程语言的 SDK 开发包,如果你是 Java 开发者,推荐使用 [Spring AMQP](https://www.rabbitmq.com/tutorials/tutorial-one-spring-amqp)。 > AMQP 是高级消息队列协议(Advanced Message Queuing Protocol)的缩写,是一个开放标准,不绑定特定技术。RabbitMQ 就是基于 AMQP 实现的,所以客户端库叫 Spring AMQP。 只需几行代码,就能创建队列、让生产者发送消息、让消费者接收处理消息: ![](https://pic.code-nav.cn/post_picture/1601072287388278786/7XJS4yWTX3qkuYy2.webp) 你:不是吧,这么简单?! 鱼皮:没错,但刚刚我们只是完成了最简单的 1 对 1 发送和接收消息。 实际上 RabbitMQ 有 6 种工作模式,适用于不同的业务场景,这也是消息队列的学习重点。 ## 第三阶段:RabbitMQ 工作模式 #### 简单模式 最简单的是 **Simple 模式**,一个生产者、一个队列、一个消费者。 就像老板派发任务给员工,**队列(Queue)** 是存储任务的容器,老板把任务放进去,员工从里面取出来完成。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/Ryb8wUNSbxU7fJMk.webp) #### 工作队列模式 鱼皮:如果有很多任务要处理,一个员工忙不过来怎么办? 你:多找几个员工帮忙? 鱼皮:没错,这就是 **Work Queue 工作队列模式**,一个生产者、一个队列、多个消费者。 就像老板发布了一堆任务,RabbitMQ 会把任务依次分配给员工,但是一个任务只会被一个员工完成。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/XybAyel26enqKXp7.webp) #### 发布订阅模式 你:这样效率就高多了! 但如果老板要求所有员工都写工作总结,怎么把同样的任务发给多个员工呢? ![](https://pic.code-nav.cn/post_picture/1601072287388278786/7iYD1PB2HMJWv6Lp.webp) 鱼皮:好问题!之前工作队列是多个员工共享一个任务列表,而现在每个员工都要有自己的任务队列。 老板需要利用 **交换机(Exchange)** 来控制把任务分发给哪些和它绑定的队列。 比如想把任务发给所有员工,就要用到 **广播交换机(Fanout Exchange)**,它会把任务发给所有已绑定的队列,然后每个员工分别从自己的队列取任务并完成。这就是 **发布订阅模式(Publish/Subscribe)**。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/lqhgm2dKrRgPftsO.webp) #### 路由模式 你:那如果老板想把某些任务发给特定员工呢?比如鱼皮负责写代码和修 Bug,小阿巴负责写代码和摸鱼~ 鱼皮:可以使用 **路由模式(Routing)**,给每个员工的队列设置自己负责的路由键(Routing Key): - 鱼皮的队列绑定 “写代码” 和 “修 Bug” - 小阿巴的队列绑定 “写代码” 和 “摸鱼” 老板发布任务时会指定路由键,由 **直接交换机(Direct Exchange)** 根据路由键精确匹配,把任务发给对应的队列。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/ghxDgMASnl4gBHDR.webp) #### 主题模式 你:那如果老板想把所有前端相关的任务都发给前端员工,后端相关的任务都发给后端员工呢?一个一个指定路由键是不是太麻烦了? 鱼皮:可以使用 **主题模式(Topic)**,它使用 **主题交换机(Topic Exchange)**,支持使用通配符模糊匹配。 可以给队列绑定通配符路由键,就能接收所有匹配的任务: - 前端员工的队列绑定 `frontend.*`(匹配 1 个词),能匹配 `frontend.Vue`、`frontend.React` - 后端员工的队列绑定 `backend.#`(匹配多个词),能匹配 `backend`、`backend.Java`、`backend.Java.优化` ![](https://pic.code-nav.cn/post_picture/1601072287388278786/6wc85ANlerLYpqVv.webp) 你:哇,这样一来就灵活多了! 鱼皮:没错,这 5 种模式是企业中最常用的,掌握了它们就能应对大部分场景了。至于最后一种模式 —— RPC(远程过程调用)模式,暂时不需要了解,因为企业开发中一般用专门的 RPC 框架,比如 gRPC、Dubbo。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/DuLmoI6M5Vs5Urym.webp) 你:好嘞,我先用这些工作模式来重构秒杀抢购功能,请好吧您嘞! ## 第四阶段:消息队列生产实践 一个月后,你意气风发地找到鱼皮:鱼皮 gie,我的抢购代码重构完了,快帮我看看能不能上线~ 鱼皮看着你的代码,表情难受得像持矢一样:emmm,你这代码要是上线,公司就完蛋了! ![](https://pic.code-nav.cn/post_picture/1601072287388278786/PL8IAgFjKCQDEltt.webp) 你:阿巴阿巴,我本地测试过了,完全没问题啊! 鱼皮:消息队列在生产环境中的坑可多着呢!下面我来问你几个问题。 #### 第一问:消息会丢吗? 鱼皮:如果重启 RabbitMQ 服务器,队列里的消息会丢吗? 你支支吾吾:会……会丢吧? 鱼皮:当然会!在抢购系统中,如果消息丢了,订单永远不会被创建,用户那边会一直显示 “抢购中”。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/7EH6jfvQWBRaXqOg.webp) 你着急了:那咋办啊,我会吃投诉的! 鱼皮:要保证消息不丢失,需要做好持久化 3 件套,把数据从内存保存到硬盘: - 队列持久化,创建队列时设置 durable 为 true。 - 消息持久化,发送消息时设置消息为持久化模式(比如 deliveryMode = 2 或 persistent = true) - 交换机持久化,创建交换机时设置 durable 为 true。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/8VAmrRUcXZwVzAdq.webp) 你:哈哈,这下消息就不会丢了吧! 鱼皮:不一定!光持久化还不够,还要有消息确认机制来保证消息的可靠性。 1)生产者确认(Publisher Confirm):生产者发送消息后,等待 RabbitMQ 的确认回复,确保消息真的被接收了。就像寄快递,你把包裹交给快递员,快递员要给你一个回执单,证明他收到了。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/r7ZRReFnJP1ZKEau.webp) 2)消费者确认(Consumer ACK):消费者处理消息成功后,要手动告诉 RabbitMQ “我处理完了,你可以删了”。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/iWoZOWEstBWEJD2j.webp) 千万别用默认的自动确认,那样消息一收到就删除,万一处理失败消息就丢了。就像收快递要签收,确保包裹真的送到你手上了。 鱼皮:这样从生产到消费的整个链路都有保障,消息就不会丢失了。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/K6tnbOig3jHBseWZ.webp) #### 第二问:消息会重复吗? 鱼皮:如果网络抖动,或者消费者重启,同一条消息可能被消费多次。比如抢购消息被重复消费,同一个用户的订单被创建了 2 次,库存被扣了 2 次怎么办? ![](https://pic.code-nav.cn/post_picture/1601072287388278786/kH7qg3O9c1zA4fNL.webp) 你:那我得跑路了! 鱼皮:不至于不至于,我们要保证 **消息幂等性**,让重复消费等同于只消费一次。 常见的方案有 3 种: 1)给每条消息一个唯一 ID(比如 UUID 或雪花算法 ID),消费前先检查这个 ID 是否处理过。 2)利用数据库唯一索引,比如订单号设置为唯一索引,重复插入会失败。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/IFWP2UysIar831Mj.webp) 3)使用 Redis 分布式锁,同一条消息同一时间只能有一个消费者处理。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/EdJapDotSRAo4hAz.webp) 你:明白了!抢购时为了防止重复下单,要检查用户是否已经下过单了。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/e4TVxH5uvfkKmY1j.webp) 鱼皮点头:没错,幂等性设计是分布式系统的基本功。 #### 第三问:消息乱序怎么办? 鱼皮:比如用户先抢购下单、然后取消订单,分别向队列发了 “创建订单” 和 “取消订单” 2 条消息。如果处理顺序乱了,系统先处理取消、后处理创建,订单状态就错了。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/PPK20qCrPr3NUYow.webp) 你:队列不是先进先出的吗?RabbitMQ 应该天然有序吧? 鱼皮:单队列内确实有序,但如果你开了多个消费者并发处理,可能消息 1 还在处理,消息 2 已经处理完了,顺序就乱了。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/gdSMSAkbklpBEgJQ.webp) 你:那怎么保证顺序呢? 鱼皮:在 RabbitMQ 中,要严格保证顺序,建议用 **单队列 + 单消费者**。 你:那性能不就很低了吗? 鱼皮:没错,所以要根据业务需求权衡一致性和性能。 如果你全都要,可以考虑 Kafka 的分区机制,可以将同一用户的消息路由到同一个分区,同一个分区内的消息严格有序,不同用户的消息又可以并发处理。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/ZrHGQpsBdju2daUH.webp) 你:妙啊,又学到一个架构知识! #### 第四问:消息处理失败怎么办? 鱼皮:如果消费者处理消息时出错了,比如数据库连接失败、业务逻辑异常,怎么办呢? ![](https://pic.code-nav.cn/post_picture/1601072287388278786/miF9p3tUVU8ypxNR.webp) 你:重试? 鱼皮:重试几次还是失败呢? 你:阿巴阿巴,直接删掉消息? 鱼皮:删掉不就丢了吗?!这时应该要用死信队列(Dead Letter Queue)。 3 种情况会产生死信: 1. 消息被消费者拒绝(reject 或 nack,requeue 参数设置为 false) 2. 消息过期 3. 队列满了 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/hwRHXJe1g1kvquHn.webp) 你:呜呜呜,这些死信好可怜。 鱼皮:没错,我们要给死信一个去处。可以配置死信交换机(DLX),将死信自动路由到死信队列。由专门的消费者监控死信队列,发现死信后,可以重试、告警、或者人工处理。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/BPmAnsRTgUfBjWPh.webp) 你:好耶,这样异常消息就不会丢失啦! #### 第五问:RabbitMQ 服务挂了怎么办? 鱼皮:前面说了这么多保证消息不丢的机制,但如果 RabbitMQ 服务本身挂了,整个系统就不能收发消息了,这可是大事故! 你汗流浃背了:那怎么办? 这时,练习两年的实习生阿坤突然鸡叫起来:我知道,要搭建集群和高可用架构! ![](https://pic.code-nav.cn/post_picture/1601072287388278786/Yon6V8E3d9KIYLIs.webp) 生产环境至少要搭建 3 个 RabbitMQ 节点的集群。 对于重要的队列,要使用 **Quorum 仲裁队列**,它基于 Raft 协议实现,会把数据自动同步到多个节点,保证数据一致性。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/DLnHT3BD01ax9xDt.webp) 如果主节点挂了,Raft 会自动选举新的主节点,实现故障转移,用户毫无感知。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/HbGW2TxIUwBe7yNy.webp) 如果数据量太大,单个队列存不下,可以使用 Sharding 插件创建 **分片队列**,把消息分散到多个节点。不过 RabbitMQ 并不擅长这个场景,建议使用 Kafka。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/jCjRPlIe0fR0uMS2.webp) 鱼皮拍了拍阿坤的肩膀:不错!生产环境的高可用是非常重要的。 此外,RabbitMQ 还有一些高级特性,在特定场景下很有用。 #### RabbitMQ 高级特性 1)延迟队列 比如想要自动取消超过 15 分钟还未支付的订单,可以用死信队列配合生存时间(TTL)实现。发消息时设置 TTL 为 15 分钟,等它过期变成死信,自动进入死信队列,然后让死信队列的消费者处理取消逻辑。但是更推荐使用专门的延迟插件 `rabbitmq_delayed_message_exchange`。 2)优先级队列 比如想优先处理 VIP 用户的订单,可以设置队列的最大优先级(x-max-priority),然后在发送消息时指定优先级,优先级高的消息会被优先消费。 3)Stream 流式存储 传统队列消息消费后就删除了,而 Stream 可以保留历史消息、重复消费、回溯历史,类似 Kafka,适用于实时数据分析、审计日志等场景。但我建议了解即可,不如用 Kafka。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/nDRpoU64CHxq4gEO.webp) 你羞愧地低下了头:我以为自己已经掌握了 RabbitMQ,原来只是学了个皮毛... ## 第五阶段:深入原理 被鱼皮连环拷问后,你主动找到阿坤:坤哥你好强,我想深入学习 RabbitMQ 底层原理,请问你是咋学的啊? 阿坤有些惊讶:咦?你不背八股文的么?[刷刷题](https://www.mianshiya.com/) 就好了呀! ![](https://pic.code-nav.cn/post_picture/1601072287388278786/MY1or8ZczwOomj1p.webp) 你震惊了:现在的校招生,竟然恐怖如斯! 鱼皮:阿坤你别逗他了,其实可以带着问题学习,重点探索 RabbitMQ 的消息路由机制、队列存储结构、AMQP 协议和消息持久化的实现。 比如: 1. 消息路由机制:Exchange 用什么算法匹配?Binding 如何存储? 2. 队列存储结构:消息在内存还是磁盘?不同队列类型有什么区别? 3. AMQP 协议:客户端和服务端怎么通信?数据格式是什么样的? 4. 持久化实现:数据怎么写入磁盘?如何保证不丢失? ![](https://pic.code-nav.cn/post_picture/1601072287388278786/OdYGIGMCiIkvvUUs.webp) 你好奇:那要怎么学习这些原理呢? 鱼皮:从这些问题出发,去阅读相关的文章,推荐官方文档和官方博客。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/wwyC42dZROCjCMoU.webp) 或者像阿坤说的刷一刷 [RabbitMQ 面试题](https://www.mianshiya.com/),就能快速学会很多核心知识点。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/d5E5rkaGYj6Ku40r.webp) 如果想增加求职竞争力,最重要的是做实战项目,我在 [编程导航](https://www.codefather.cn/) 上的的智能 BI 项目和 OJ 判题系统项目都有完整的 RabbitMQ 实战。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/2RJ4lgUn1nHf41gY.webp) 你握紧了粉拳:好的,我这就去学! ## 结尾 若干年后,你已经成为了大厂的消息队列专家。不仅能熟练使用 RabbitMQ 解决各种业务问题,搭个集群架构也是手拿把掐的。 你也像鱼皮当时一样,耐心地给新人分享学习 RabbitMQ 的经验,让他们谨记 “消息队列是实战型技术,一定要多动手实践”。 ![7a5e7627-2df3-4b32-ad2f-3231ade8015d.png](https://pic.code-nav.cn/post_picture/1601072287388278786/O0LGjxL4U1MlCQN8.webp) 再次遇到鱼皮是在一条昏暗的小巷,此时的他年过 35,一毛不拔。你什么都没说,只是给他点了个赞。 不打扰,是你的温柔。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/YOoA3yZdnae71cmE.webp) ## 更多 💻 编程学习交流:[编程导航](https://www.codefather.cn/) 📃 简历快速制作:[老鱼简历](https://www.laoyujianli.com) ✏️ 面试刷题神器:[面试鸭](https://www.mianshiya.com) 📖 AI 学习指南:[AI 知识库](https://ai.codefather.cn/)

消息队列

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

MQ进阶 - 分布式下的挑战

上文讲了三大主流 MQ 的的功能和特性,那就像是工具箱,在实际场景中按需使用即可。 本文讲述的是使用 MQ 我们会遭遇的问题,相当于 BUG —— 使用 MQ 出现了不符合预期的结果,不解决不可上线。 ## 引言 - 分布式的经典问题 在分布式场景下,给我们带来了高可用、易扩展和高容错的便利,也带来了分布式系统下不少的麻烦,增加了分布式开发的复杂性。 分布式的问题一图概览: ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/UkJfS0f7ooublNIb.svg) 可知主要问题是三个根源引起:网络不可靠、节点独立性和无全局时钟。 **1)网络不可靠**:在我们物理世界,网络波动(**延迟、丢包、分区甚至中断**)是无法根除的情况,那么这会导致: - **消息丢失或重复** - **网络分区**: 分布式系统中,不同节点之间因为网络故障,**无法互相通信**,从而被“分隔成多个孤岛”。 某个节点故障必然形成网络分区。 **2)节点独立性**:分布式系统由多个独立计算机(节点)组成,每个节点都有自己的内存、CPU和时钟。独立节点可能会出现: - **节点故障**:任何节点都可能随时发生故障(宕机、网络中断、磁盘损坏等),这与单机程序中“要么全好、要么全坏”的模式截然不同,系统必须能处理**部分失效**(Partial Failure)的情况。 3)无全局时钟:虽然每个节点都有自己的物理时钟(如石英钟),但它们之间**无法做到精确同步。** - **时钟不同步**:即使通过NTP(网络时间协议)同步,也存在毫秒级甚至更大的误差,且时钟会因温度、电压等因素漂移 - **事件顺序难定**:因为时钟不同步,节点A记录的事件时间戳可能早于节点B,但实际发生顺序却相反。这使得**确定跨节点事件的绝对顺序变得不可能** **为什么大家使用北京时间,还会出现时间不同步呢?** 服务器通常会定期通过 NTP 同步时间到“标准北京时间”(如国家授时中心)。 但是同步不是实时的,比如每隔几分钟或几小时才同步一次。 而网络延迟、抖动、负载都会导致同步结果存在 误差(几十毫秒到几百毫秒)。 另外 每台机器内部都有一个石英振荡器来计时。 操作系统每隔一段时间会读这个振荡器的值来更新系统时间。 但是!不同机器的振荡器频率略有差异(例如每秒快几微秒、慢几微秒),时间会逐渐漂移(drift)。 以上都是在分布式系统下无法突破的物理层面限制,这些物理层面的问题,我们必须在软件层面克服,物理层面无法突破,因此下列问题是软件开发者必须面对的经典挑战: | 问题类型 | 具体表现 | 根源关联 | | ---------------------------- | ------------------------------------------------------------ | -------------------------------------- | | **数据一致性** | 数据在多个副本间可能出现不一致。 | 网络分区、节点故障、消息延迟与重复。 | | **可用性** | 某节点失联或网络中断后,服务完全中断或响应严重变慢 | 节点独立性、网络不可靠 | | **分布式共识** | 多个节点如何对某个值(如谁是Leader)达成一致 。 | 节点独立性、网络不可靠。 | | **分布式事务** | 确保跨多个节点的操作序列要么全部成功,要么全部失败 ,分布式难以实现 | 网络不可靠 、节点故障、无全局时钟。 | | **故障检测与恢复** | 如何准确判断一个节点是真的宕机了还是只是网络慢,以及如何恢复。 | 网络不可靠(延迟与分区)、节点独立性。 | | **消息可靠性(丢失/重复)** | 消息在网络中丢失、延迟或被重复投递,导致消费者状态错误 | 网络不可靠、节点独立性 | 将分布式的挑战讲的深入一点,是为了看清楚消息队列的重复消费、幂等性、一致性问题是由分布式架构带来的,根源是出自物理层面不可避免的问题,我们来学习前辈如何巧妙的在软件层面攻破物理层面的困难。 | 问题类型 | 具体表现 | 常见解决方案 | RabbitMQ 支持情况 | RocketMQ 支持情况 | Kafka 支持情况 | | ---------- | ------------------------------------------------------------ | ------------------------------------------------------- | ------------------------------------------------------------ | ------------------------------------------------------------ | ------------------------------------------------------------ | | 消息丢失 | 生产者发送失败、Broker 存储失败、消费者处理失败但未确认,导致消息永久丢失。 | 生产者确认机制、Broker 持久化、消费者手动确认、副本机制 | ✅(发布者确认、持久化队列/消息、手动确认)[rabbitmq.com](https://www.rabbitmq.com/docs/reliability?utm_source=chatgpt.com)[rabbitmq.com](https://www.rabbitmq.com/tutorials/tutorial-seven-java?utm_source=chatgpt.com) | ✅(同步/异步发送、刷盘策略、主从复制)[AlibabaCloud](https://www.alibabacloud.com/blog/how-rocketmq-helps-achieve-better-message-reliability_597836?utm_source=chatgpt.com)[Stackademic](https://blog.stackademic.com/rocketmq-practical-guide-a-complete-analysis-of-distributed-messaging-middleware-for-java-6283563e1ee5?utm_source=chatgpt.com) | ✅(acks = –1(所有副本确认)、副本机制、消费者位移提交) | | 重复消费 | 网络延迟或确认失败导致消息被重投,或消费者重启后重新消费,导致同一消息被处理多次。 | 消费端幂等设计、唯一ID去重(DB/Redis) | ❌(主要依赖业务设计) | ❌(推荐业务唯一标识) | ✅(生产者幂等(单分区/单会话)、流处理 Exactly Once 支持;但消费端仍需幂等处理) | | 消息顺序 | 需要保证一组有序的消息(如同一订单的创建→支付→发货)按照发送顺序被消费。 | 全局有序(单队列)或分区有序(按 key 路由同一队列) | ✅(队列内 FIFO;全局有序需特殊配置) | ✅(分区有序:顺序消息支持相同 Key 到同一队列) | ✅(分区内有序:相同 Key 发到同一分区) | | 消息积压 | 生产者速度远超消费者,或消费者故障,导致消息在 Broker 中大量堆积。 | 水平扩展(增加消费者/分区)、监控告警、限流降级 | ✅(监控队列长度、可增加消费者) | ✅(高性能堆积场景支持、监控与扩展能力强) | ✅(天然支持海量堆积、可增加分区和消费者) | | 事务消息 | 跨多个系统/数据库的操作必须具有原子性(例如订单创建成功才发送消息通知下游)。 | 两阶段提交(2PC)、本地消息表、事务消息机制 | ❌(原生不支持,仅可通过插件或本地消息表+Confirm 模拟) | ✅(原生支持事务消息机制:半消息 → 提交/回滚) | ✅(支持流处理场景的事务/Exactly Once,但并非传统队列事务机制) | | 分布式事务 | 涉及多个消息生产与消费的完整业务流程需保证最终一致性。 | Saga 模式、TCC 模式、可靠消息最终一致性 | ❌(主要依赖业务模式:Saga 或本地消息表结合) | ✅(提供基础机制:事务消息为可靠消息最终一致性利器) | ✅(提供基础机制:在流处理里支持事务(EOS)) | | 重试补偿 | 消费失败后自动/手动重试,直到达到上限或人工干预(补偿)处理失败消息。 | 自动重试机制、死信队列、人工补偿 | ✅(**原生支持死信队列**,消息被拒绝或过期时可路由至DLX) | ✅(**原生支持消费重试**,达到上限后进入**死信队列**) | ❌(**无内置重试机制**,需业务方自行实现,如发到重试Topic) | - **事务消息**:主要针对的是“**本地事务操作与消息发布之间的一致性问题**”,也就是当你做操作(例如数据库写入)然后发送消息通知下游时,如何做到“要么操作+消息都成功,要么都失败”。 - **分布式消息**:泛指在分布式系统中各节点之间通过消息通信所涉及的“异步、跨节点、解耦、可靠性”问题。它更宽泛,不仅包括事务相关,还包括消息传递、失败重试、顺序、分区、幂等等。 在进入这些问题的学习之前,我们再了解一个分布式下的必要知识点,CAP定理。 ## CAP 定理 CAP : 一致性(Consistency)、可用性(Availability)、分区容错性(Partition tolerance) 根据[百度百科](https://baike.baidu.com/item/CAP原则/5712863?fromtitle=CAP定理&fromid=22739568): CAP定理的思想由**加州大学伯克利分校的计算机科学家埃里克·布鲁尔** 于**2000年**在分布式计算原则研讨会(PODC)上首次提出,当时被称为**CAP猜想。** **2002年**,**麻省理工学院(MIT)的赛斯·吉尔伯特和南希·林奇** 发表了正式证明,从理论上验证了这一猜想的正确性,使其成为一个**定理。** ### CAP是什么? 在分布式系统健康的状态下,不存在CAP理论,因为分区容错性,一致性和可用性完全可以得到保障。当单个节点发生故障时,必须优先保证分区容错性,此时 CAP 定理发挥作用。 CAP定理指出,任何一个**分布式系统**(由多个通过网络通信的节点组成的系统)**无法同时同时满足**以下三个基本需求,最多只能同时满足其中的两个 - **一致性**: 所有节点在同一时间看到的数据是一致的。 - **可用性**: 每个请求都会得到响应(不管成功或失败)。 - **分区容错性**: 系统能继续运行,即使网络部分节点无法通信。 CAP 定理揭示了**分布式系统中的一种固有状况和必然的权衡,**一种无法解决的困境**。** ### 为什么不能三个都满足 理想情况下,我们肯定希望三个都满足,如果不能,我们必须要优先满足分区容错性,因为网络故障是不可避免的,如果分区容错性不能保证,那么其它两个更加做不到。 如果选择**保证一致性(C)**:那么在分区期间,某些节点可能需要停止服务(拒绝读写请求),以避免数据不一致,这就**牺牲了可用性(A)** 如果选择**保证可用性(A)**:那么系统即使在分区期间也必须响应每个请求,但这可能意味着返回的是旧数据(不同节点数据可能不一致),这就**牺牲了一致性(C)** ### 举例论证 系统故障,发生以下网络分区的情况: - 假设系统有三个节点 A、B、C - 节点 A 和 B 可以互相通信。 - 节点 C 与 A/B 网络断开(分区)。 - 现在系统出现了**孤岛**:A/B 和 C 之间无法通信。 如果**选择一致性**,可用性必然下降: - 系统必须保证每次写入都被所有节点确认(或者按某种算法保证一致性)。 - 但是节点 C 无法与 A/B 通信。 - 如果继续让 C 提供写入或读取,就可能看到不同步的数据 → 破坏一致性。 - 所以为了 **保证一致性**,节点 C 就必须 **拒绝读写请求**,相当于停服。 如果**选择可用性,**一致性必然下降: - 即使 C 与 A/B 失去连接,C 仍然响应请求(读/写)。 - 但是 C 的数据可能与 A/B 不一致(写入只在 C,未同步到 A/B)。 - 这就破坏了全局一致性(C)。 强行让 C 服务器运转,造成数据不一致。 ## 消息队列的基本流程 我们先看看 RocketMQ 和 RabbitMQ 生产消息和消费消息的最基本流程。 ### RocketMQ 核心角色 - **Producer (生产者)**: 消息的发送方。负责产生消息,并发送到指定的 Topic。 - **Consumer (消费者)**: 消息的接收方。负责从指定的 Topic 订阅并消费消息。 - **Broker (代理服务器)**: 消息中转角色。负责存储来自生产者的消息,并将消息投递给消费者。是消息存储和转发的核心。 - **NameServer (注册中心)**: 无状态的轻量级服务。充当服务发现的角色,维护了 Broker、Topic 等路由信息。生产者和消费者会从 NameServer 获取目标 Broker 的地址列表。 **RocketMQ 生产消息到 Brocker 的详细流程:** 1. **创建生产者**: 应用程序创建一个 `**DefaultMQProducer**` 实例,并设置 NameServer 的地址。 2. **启动生产者**: 调用 `**producer.start()**` 方法。 - - 内部会启动多个服务线程。 - 与 NameServer 建立长连接,定期(默认每30秒)从 NameServer 拉取最新的 Topic 路由信息(即该 Topic 对应的 Broker 地址、队列等)。 1. **构建消息**: 创建一个 `**Message**` 对象,包含 Topic、Body(消息体)和可选的 Tag/Keys。 2. **发送消息**: 调用 `**producer.send(message)**` 方法。 - - **选择队列**: 生产者客户端会根据预设 (代码逻辑预设) 的负载均衡策略(如轮询),从获取到的路由信息中选择一个消息队列(Message Queue,位于某个 Broker 上)。 - **网络发送**: 将消息发送给目标 Broker。 1. **等待 Broker ACK (关键步骤)**: - - 默认情况下,`**send()**` 方法是**同步阻塞**的。它会一直等待,直到收到来自 Broker 的响应,或者直到超时。 - Broker 接收到消息后,会将其写入 CommitLog(内存和磁盘),然后向生产者返回一个 `**SendResult**` 对象。 - 这个 `**SendResult**` 就是**生产端的 ACK**。它包含了发送状态(`**SendStatus**`),如 `**SEND_OK**`(成功)、`**FLUSH_DISK_TIMEOUT**`(刷盘超时)、`**FLUSH_SLAVE_TIMEOUT**`(同步到从机超时)、`**SLAVE_NOT_AVAILABLE**`(从机不可用)等。 1. **处理结果**: 生产者根据 `**SendResult**` 的状态,可以判断消息是否成功发送,并进行相应的业务处理(如重试、记录日志等)。 ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/wlfrHSge2RgUmaGP.svg) **RocketMQ 生产消息到 Brocker 的详细流程:** 1. **创建消费者**: 应用程序创建一个 `**DefaultMQPushConsumer**` 实例,设置 NameServer 地址、消费者组名(Consumer Group)。 2. **订阅 Topic**: 调用 `**consumer.subscribe("TopicName", "\*")**` 方法,订阅感兴趣的 Topic。 3. **注册消息监听器**: 调用 `**consumer.registerMessageListener()**` 注册一个监听器实现(例如 `**MessageListenerConcurrently**`)。这个监听器包含了**业务处理逻辑**。 4. **启动消费者**: 调用 `**consumer.start()**` 方法。 - - 内部会与 NameServer 建立长连接,获取 Topic 路由信息和 Broker 地址。 - 根据消费者组名和负载均衡策略,为当前消费者分配若干个消息队列进行消费。 - 启动后台拉取线程,**主动**向其负责的 Broker 拉取消息。 1. **Broker 投递消息**: Broker 的拉取处理器接收到消费者的请求后,从存储中读取消息并返回给消费者,如果没有消息会进行长轮询。 2. **消费消息 (业务处理)**: - - 消费者客户端收到消息后,会调用用户注册的 `**MessageListener**` 的 `**consumeMessage**` 方法,将消息交给业务代码处理。 1. **返回消费状态 (关键步骤)**: - - 业务代码处理完毕后,`**consumeMessage**` 方法必须返回一个 `**ConsumeConcurrentlyStatus**` 状态。 - `**CONSUME_SUCCESS**`: 表示消费成功。 - `**RECONSUME_LATER**`: 表示消费失败,希望稍后重新消费。 1. **发送消费端 ACK (关键步骤)**: - - **如果** `**consumeMessage**` **返回** `**CONSUME_SUCCESS**`:消费者客户端会**自动**向 Broker 发送一个 ACK 请求。 - **如果** `**consumeMessage**` **返回** `**RECONSUME_LATER**` **或抛出异常**:消费者客户端**不会**发送 ACK。 1. **Broker 处理 ACK**: - - Broker 收到 ACK 后,会将该消息在消费队列中的偏移量更新,标记为已消费。 - 如果在一定时间内(默认约15分钟)未收到某条消息的 ACK,Broker 会认为消费失败,并重新将该消息投递给消费者组内的其他消费者(如果有的话),或在下次拉取时再次投递给原消费者。 ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/hD3HLmUXPWQfyG4y.svg) ### RabbitMQ 核心角色 - **Producer (生产者)**: 消息的发送方。负责将消息发送到交换机。 - **Consumer (消费者)**: 消息的接收方。负责从队列中获取消息并进行处理。 - **Broker (代理服务器)**: RabbitMQ 服务本身。负责接收、存储和路由消息。 - **Exchange (交换机)**: 接收来自生产者的消息,并根据路由键将消息推送到一个或多个队列。生产者从不直接将消息发送到队列。 - **Queue (队列)**: 存储消息的缓冲区,直到消费者将其取走。 - **Channel (信道)**: 位于 TCP 连接内部的虚拟连接。它是进行 AMQP 操作(如发布、消费、获取消息)的轻量级载体,避免了为每个操作都建立 TCP 连接的开销。 **RabbitMQ 生产消息到队列的详细流程:** 1. **创建连接和信道**: 应用程序与 Broker 建立一个 TCP 连接,并在其上创建一个信道。 2. **启用发布者确认 (关键步骤)**: 在信道上调用 `**channel.confirmSelect()**`,将该信道设置为确认模式。**这是开启生产端 ACK 的前提**。 3. **声明交换机和队列**: 为了保证消息能被正确路由,生产者通常会声明一个交换机和一个队列,并将它们绑定起来(这是一个幂等操作,重复声明不会出错)。 4. **发送消息**: 调用 `**channel.basicPublish()**` 方法,将消息发送到指定的交换机,并附带路由键。 5. **等待 Broker ACK (关键步骤)**: - - 消息发送后,生产者不会立即收到响应。Broker 接收到消息后,会将其路由到匹配的队列。 - 一旦消息至少被一个队列接收并持久化(如果队列是持久化的),Broker 就会通过该信道向生产者发送一个**确认帧**,包含消息的 `**delivery-tag**`。 - 如果由于某些原因(例如找不到匹配的队列)消息无法被路由,Broker 会发送一个**未确认帧**。 1. **处理确认**: 生产者可以监听信道的确认和未确认事件,根据收到的 `**delivery-tag**` 来判断哪条消息成功或失败,并进行相应的补偿业务处理(如记录日志、重发失败的消息等)。 ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/lYUabk1qaVW4ce21.svg) **RabbitMQ 消费流程:** 1. **创建连接和信道**: 与生产者类似,消费者也需要建立连接和信道。 2. **声明队列**: 消费者声明它要消费的队列,确保队列存在。 3. **设置 QoS (可选但推荐)**: 调用 `**channel.basicQos(prefetchCount)**`,限制消费者每次能从队列预取的消息数量。这可以防止消费者被消息淹没,实现更公平的负载均衡。 4. **定义消费者并消费**: 创建一个 `**Consumer**` 实现类(或使用回调),并调用 `**channel.basicConsume()**` 方法开始消费。 - - **关键参数**: 在 `**basicConsume**` 中,必须设置 `**autoAck**` 为 `**false**`。这表示**手动确认模式**。如果设置为 `**true**`(自动确认),RabbitMQ 会在消息发送给消费者后立即删除它,不管消费者是否处理成功,这是极不可靠的。 1. **Broker 投递消息**: RabbitMQ 将队列中的消息推送给消费者。 2. **消费消息 (业务处理)**: 消费者的回调方法被触发,消息内容被传递给业务代码进行处理。 3. **发送消费端 ACK (关键步骤)**: - - **如果业务处理成功**: 在 `**handleDelivery**` 回调方法的最后,调用 `**channel.basicAck(deliveryTag, multiple)**`。`**deliveryTag**` 是消息的唯一标识,`**multiple**` 表示是否确认小于该 `**deliveryTag**` 的所有消息。 - **如果业务处理失败**: 调用 `**channel.basicNack(deliveryTag, multiple, requeue)**` 或 `**channel.basicReject(deliveryTag, requeue)**`。`**requeue**` 参数决定消息是否重新返回队列头部等待再次投递。 1. **Broker 处理 ACK**: - - Broker 收到 `**basicAck**` 后,才会从队列中移除该消息, 会删除持久的消息。 - 如果收到 `**basicNack**`/`**basicReject**` 且 `**requeue=true**`,消息会被重新放回队列。 - 如果消费者在处理消息时断开连接且未发送 ACK,Broker 也会认为消费失败,并将消息重新入队。 ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/iHZsU0ytuwkPOUcN.svg) ## 消息丢失 **消息丢失**指的是:生产者发送了消息,但这条消息**未能被成功存储或被消费者处理**,最终在系统中“消失”,不再被传递或消费。消息丢失问题会严重影响数据一致性和业务可靠性 ### 问题产生 消息丢失通常有以下几个环节可能发生: - **生产者发送失败**:生产者向 Broker 发送消息时,可能因为网络中断、Broker 不可用、客户端超时而失败。如果生产者误以为发送成功但实际上没有,则消息丢失。 - **Broker 存储失败或确认失败**:消息到达 Broker,但在持久化(写磁盘)或写入日志、刷盘、同步副本等过程中失败,或者 Broker 挂掉,还未完成持久化。 - **消费者处理失败并且未重试/未记录**:消费者接收消息后处理失败,且系统没有做重试或者死信机制、也没有记录该消息状态,导致消息被“吞掉”。 - **配置或使用不当**:例如队列或主题未配置持久化、生产者未启用确认机制、Broker 副本机制配置不当、网络故障未被捕捉等。 总结一句话:消息丢失的根因通常是**网络、存储、确认机制不可靠**,原因很多,无法全部概括。 ### 问题解决 确保消息不丢失需要从**生产者、Broker、消费者**三个层面进行保障 1. **生产者端**: - - **使用发送确认机制**:确保消息成功发送到Broker。 - **配置重试机制**:发送失败时自动重试 1. **Broker端**: - - **消息持久化**:将消息和队列设置为持久化,写入磁盘 - **配置副本机制**:通过主从复制或多副本保证数据冗余 1. **消费者端**: - - **手动提交offset**:在消息处理完成后手动确认(程序员代码中用业务逻辑控制,不依赖消息中间件提供的自动提交),避免自动提交导致的问题 - **消费幂等性设计**:确保重复消费不会导致业务逻辑错误 但是这不能百分百保证消息丢失,所以下一个小节会讲消息补充,我们能接受极端情况下的丢失。 #### RabbitMQ 解决 RabbitMQ通过**多层次确认机制**和**持久化策略**保障,其核心思想是在消息传递的每个关键节点设置确认点。 1. 生产者确认机制 - - **Publisher Confirms:** 当生产者发送消息到RabbitMQ后,Broker会返回一个确认结果(ACK / NACK) - **Publisher Returns :** RabbitMQ 提供了 **ReturnCallback** 回调机制 ,开启**mandatory=true** 发送消息 时, 如果交换机找不到任何队列匹配,Broker 会 **立即回调生产者**,告诉你消息无法路由。 这样生产者就知道有消息没被队列接收,可以选择 **重试、记录日志或发送到死信队列**。 1. 持久化策略 为了保证消息不丢失,RabbitMQ不仅仅持久化了消息,还持久化了 Exchange 和 Queue , 做到了三层立体式持久化。 | 层级 | 持久化对象 | 含义 | 作用 | | ------------- | ---------------- | --------------------------------------------------------- | ------------------------ | | ① Exchange 层 | **交换机持久化** | Broker 重启后,交换机配置(名称、类型、绑定关系)仍然存在 | 确保消息路由基础结构不丢 | | ② Queue 层 | **队列持久化** | Broker 重启后,队列仍然存在(包括绑定信息) | 确保消息能有地方落地 | | ③ Message 层 | **消息持久化** | 消息内容写入磁盘(而不是只在内存) | 确保消息数据不丢失 | 1. 消费者确认机制 消息确认是消费者向Broker发送确认,表示这条消息处理完了, Broker 收到 ack 后,才会把这条消息**从队列里删除。** **RabbitMQ 支持自动确认,**缺点是容易丢失消息,因为自动确认机制是在消费者真实消费数据之前发送 ACK 。自动确认吞吐量高。 **手动确认**:可以自己写代码,灵活控制发送ACK的时机,更稳定。 #### RocketMQ 解决 1. 生产者可靠性机制 - **同步发送**:使用同步发送方式,等待Broker返回确认结果,例如`**SendResult sendResult = producer.send(msg)**`。通过检查`**sendResult.getSendStatus()**`判断是否发送成功,失败时可进行重试 - **事务消息**:对一致性要求极高的场景,使用事务消息。它通过两阶段提交(2PC)确保本地事务与消息发送的原子性 1. Broker端可靠性机制 - **同步刷盘(SYNC_FLUSH)**:消息写入内存后立即刷盘,确保数据落盘,防止因Broker突然宕机导致内存中未刷盘的消息丢失 - **异步刷盘(ASYNC_FLUSH)**:消息写入内存后异步刷盘,性能更好 - **主从复制**:配置主从Broker。将`**brokerRole**`设置为`**SYNC_MASTER**`(同步复制)。Master接收到消息后,会**同步等待Slave复制成功**,才向生产者返回ACK。即使Master宕机,数据也在Slave上有备份 1. 消费者可靠性机制 - **消费成功后才提交Offset:**消费者在**业务逻辑执行成功后**,再向Broker提交消费进度。如果处理过程中消费者宕机,消息会重新投递,不会丢失 - **避免异步消费**:RocketMQ官方建议不要在消费端使用异步处理消息,以免因未处理完就提前提交Offset导致消息丢失 ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/IJm2UCtSBbRevpWV.svg) 消息丢失问题不能根除(ack丢失,刷盘崩溃等), 接受极端情况丢失,但用补偿机制兜底 ## 消息补偿 消息补偿是一种**事后补救机制**,用于在消息传递或处理失败时,通过重新发送或执行补偿操作,确保业务数据最终达到一致状态,它通常用于处理**分布式事务**中的异常情况,特别是在采用**最终一致性**模型的系统中。 具体来说,生产者发送消息失败、消费者消费失败、消息超时未确认、事务状态未知会触发消息补偿。 ### rabbitMQ 消息补偿机制 1. 生产者确认机制 生产者发送消息时,可开启 **Confirm 模式**。Broker 会异步确认消息是否到达交换机,会出现两种情况: - **ACK**:消息成功到达交换机。 - **NACK**:消息未到达交换机。生产者可在回调中捕获 NACK,执行**重新发送**或**记录日志**等补偿操作 1. 死信队列 当消息成为“死信”(如消费者拒绝且不重新入队、消息过期、队列满)时,若配置了死信交换机(DLX),消息会被路由到死信队列。进入死信队列的消息,可在程序中写逻辑对死信队列进行重试、重发(让生成者再发一次)、记录日志、人工介入。 1. 消息重试与消费者确认 消费者处理失败时,可拒绝消息并选择**不重新入队**(`**requeue=false**`),让消息进入死信队列,避免无限重试导致阻塞。同时,应确保业务逻辑的**幂等性**,以应对可能的消息重复 ### RocketMQ 消息补偿机制 RocketMQ 提供了更强大的**事务消息**机制,并依赖**定时任务**和**状态回查**来实现补偿。 1. 事务消息 **补偿机制(事务状态回查)**:若Broker长时间未收到确认,会**反向回查**生产者的事务状态,生产者需实现回查接口,根据本地事务状态决定提交或回滚。事务回查触发消息补偿。 2. 定时任务与消息重试 **生产者**:发送失败时可自动重试,自带消息补偿。 **消费者**:消费失败时,消息会自动重试。超过最大重试次数后,消息会进入**死信队列,**可监控死信队列并进行补偿。 1. 本地消息表 对于未使用事务消息的场景,可采用**本地消息表**方案 本地消息表: 在生成者端建立一个消息表,用来把消息存储在业务数据库中,和业务操作放在同一个本地事务里提交。 这样就能确保业务操作成功的同时,消息也一定成功发送,避免中间的缝隙发生其它意外导致业务操作成功,消息发送失败的情况。 本地消息表是一种补偿机制, 先**落地保存“要发的消息”**,等网络恢复、Broker可用,再**补偿性地重发**。 如果没有收到 Broker 的 ACK 确认,或者Broker 返回失败回调 ,本地消息表发挥作用。 ## 重复消费 **同一条消息被消费者处理了多次**。 消息在 Broker(消息中间件)只存在一条,但消费者因为确认或网络问题,**重复地拿到同一条消息去执行业务逻辑**。 导致重复消费的原因: - ACK 没发出去 ,Brocker 重试机制触发 - 重启后没来得及提交 offset 或确认状态,又从上次位置重新消费。 - 手动补偿/重放逻辑执行多次 , 运维或定时任务手动重发时,逻辑重复执行。 **重复消费的危害** - **数据错误**:比如重复创建订单、重复扣款,导致业务数据混乱 - **资源浪费**:无谓地消耗CPU、内存和数据库连接等系统资源。 - **业务逻辑混乱**:比如库存被多扣减,导致超卖 - **用户体验差**:用户可能收到多条重复的短信或通知 ### 解决重复消息的思路 无法避免消息不会重复发送,所以解决思路应该是**消费的幂等性**。 即无论消息被消费多少次,最终的结果都与消费一次相同。方案如下: #### 1. 唯一ID + 去重表 为每条消息生成一个**全局唯一ID**(如订单号、UUID),并在消费端创建一个**去重表**(通常使用数据库)。 - **消费前**:消费者在消费前,先查询去重表中是否存在该ID。 - **如果存在**:说明已处理,直接跳过。 - **如果不存在**:执行业务逻辑,并将该ID插入去重表 ```sql -- 去重表示例 CREATE TABLE `message_consumed` ( `message_id` VARCHAR(64) PRIMARY KEY, `consume_time` TIMESTAMP DEFAULT CURRENT_TIMESTAMP ); ``` #### 2. Redis 原子操作 利用Redis的`**SETNX**`(SET if Not eXists)命令实现原子性去重。 - **消费前**:尝试以消息ID为key向Redis写入一个值(并设置过期时间)。 - **写入成功**:说明首次消费,执行业务逻辑。 - **写入失败**:说明已消费,直接跳过 ```java // 伪代码示例 String messageId = message.getMessageId(); Boolean isSuccess = redisTemplate.opsForValue().setIfAbsent(messageId, "1", 30, TimeUnit.MINUTES); if (Boolean.TRUE.equals(isSuccess)) { // 首次消费,执行业务逻辑 processBusiness(message); } else { // 重复消息,直接跳过 logger.info("重复消息,已跳过: {}", messageId); } ``` #### 3. 数据库唯一约束 对于业务数据本身有唯一标识的场景(如订单号),可以直接在数据库表中设置**唯一索引**或**主键约束**。 即使重复消费导致重复插入,数据库会因约束冲突而报错,应用捕获异常后忽略即可 ```sql -- 订单表order_id设置唯一约束 CREATE TABLE `orders` ( `order_id` VARCHAR(64) PRIMARY KEY, -- 其他字段... ); ``` #### 4. 业务状态机判断 对于有明确状态流转的业务(如订单状态:待支付→已支付→已发货),可以通过检查当前业务状态来判断是否需要处理。 如果消息要求的操作(如“支付”)对应的当前状态(如“已支付”)已经符合,则说明是重复或过时消息,直接丢弃 #### 5. 消息添加辅助字段 在 **RocketMQ** 中,当生产者发送消息时,可以给每条消息指定一个“业务唯一键”,例如订单号、支付流水号: ```java Message msg = new Message( "OrderTopic", "TagA", "KEY_1001", // 业务唯一键 "Order created".getBytes() ); producer.send(msg); // 消费端 message.getKeys() ``` 消费者拿到 Key 后,可以在自己系统里判断该消息是否已处理过。 这个 key 是 RocketMQ 给你暴露的幂等辅助字段 ,通过它实现幂等。 RabbitMQ 消息结构更加通用,可以设置**在** `**message_id**` **或 headers 里**。 #### 6. 乐观锁 在数据表中增加一个 `**version**` 字段(版本号)。更新数据时,需要带上这个版本号。 ```sql UPDATE products SET stock = stock - 1, version = version + 1 WHERE id = 123 AND version = 5; ``` 应用在更新前先读取 `**version=5**`,执行更新时,在 `**WHERE**` 条件中带上这个版本号。然后检查受影响的行数: - 如果行数为1,说明更新成功。 - 如果行数为0,说明在你读取和更新之间,已经有其他请求修改了这条数据(version已经不是5了),本次更新失败。 通过 version 字段来确保不会重复执行,更新数据成功 version + 1 , 这比数据库字段唯一约束更宽松,但实现复杂度高了一点。 ## 消息顺序 **消息顺序**(Message Ordering)指的是消息按照生产者发送的先后顺序被消费者接收和处理。在消息队列中,这通常分为两种类型 1. **全局有序**:整个主题(Topic)中的所有消息,都严格按照生产者发送的先后顺序进行消费。例如,发送顺序是 M1, M2, M3,消费时也必须是 M1, M2, M3。 2. **分区有序(消息键分区)**:只需要保证具有某个相同特征(如同一个订单ID、同一个用户ID)的消息能够被顺序消费即可。例如,订单 `**A**` 的消息 `**A1, A2, A3**` 顺序消费,订单 `**B**` 的消息 `**B1, B2, B3**` 也顺序消费,但 `**A1**` 和 `**B1**` 谁先消费无所谓。 ### 消息顺序的重要性 操作的顺序至关重要,乱序会导致业务逻辑错误或数据不一致。 比如订单的状态流转:创建订单 → 支付订单 → 发货 → 确认收货,一个状态是一条消息,上游执行成功会给下游发送一个消息,如果消息不会顺序消费,那么会发生数据不一致,业务混乱。 在这个场景中,每一个状态肯定是顺序发给队列的,只要确保同一个订单的消息进入同一个队列,那么在队列中的消息就是物理有序,RocketMQ 默认是消费组内单实例的单线程消费,不会造成乱序的问题。 作为开发者,我们只要保证同一个订单的消息发送给固定的一个队列就好,这就要用到**消息键。** ### 消息乱序的原因 消息队列(如 RabbitMQ、Kafka、RocketMQ、Redis Stream)都是**并发消费模型**。 如果你启用了多个消费者或多线程消费,一个消息队列中不同消息可能被**不同线程同时处理**。 ### 实现消息有序 我们知道有消息顺序:全局有序和分区有序,其中分区有序是重点难点,使用的也是比较多的,因为分区有序性能高,但是复杂一点。 分区有序使用**消息键模式**实现,接下来看各个MQ怎么处理。 #### kafka 与 Rocket MQ 实现方案相同 消息键:在生成消息的时候指定 **一个能区分业务功能、能够分组业务的字段 ,**字段可以是任何类型,会经过哈希计算,对Brocker的总队列取模,相同的业务便投入到固定队列。 以kafka为语境进行讲解。 **原理详解:** 1. 生产者将消息路由到固定**分区(也就是队列)** 2. 1. **定义业务键**:为每条消息设定一个能够标识其业务分组的键。例如,处理订单消息,`**order_id**` 就是绝佳的业务键。 2. **根据键计算分区**:生产者在发送消息时,不是随机选择分区,而是**根据业务键进行哈希计算**,然后**对分区总数取模**,从而确定该消息应该发往哪个分区。 3. 示例代码: ```java // 业务消息 Message message = new Message("topic_order_paid", "订单A-已支付"); String businessKey = "order_A"; // 例如,订单ID // 计算目标分区 int partitionCount = 4; // 假设主题有4个分区 int targetPartition = Math.abs(businessKey.hashCode()) % partitionCount; // 发送到指定分区 producer.send(message, targetPartition); ``` **效果**:所有 `**order_A**` 的消息(如“已支付”、“已发货”、“已签收”)都会被哈希算法路由到同一个分区(比如分区 2)。同样,所有 `**order_B**` 的消息也会被路由到同一个分区(比如分区 0)。 1. **消费者端单线程消费分区** 1. 1. **天生支持单线程消费:**消息队列的消费者组模型天然支持这一点。**一个分区在同一时间只能被消费者组里的一个消费者实例消费。** 2. 多个消费组可以订阅一个topic, 那么多个消息组都有自己的 offset , 一条消息会被多个消费组消费,这是**跨消费组广播模式**,按需使用。 3. **组协调器:**运行在MQ服务上,管理一个消费组的元数据。消费者实例启动后回去找这个组协调器,加入进去由组协调器管理消费者,组协调器会给每个消费者分配分区列表,列表由kafka客户端维护,通过 getAssignedPartitions(); 获取,即可开始监听消费。 示例代码: ```java // 消费者启动时,会从MQ服务端分配到一些分区,例如 partition-2 // getAssignedPartitions() 返回的是分配给【当前这一个消费者实例】的分区列表,而不是整个Topic的所有分区。 List<Partition> assignedPartitions = getAssignedPartitions(); // [partition-2] // 为每个分区启动一个独立的消费线程 for (Partition partition : assignedPartitions) { new Thread(() -> { while (true) { // 从 partition-2 中拉取一条消息 Message message = consumer.poll(partition); if (message != null) { // 同步处理消息,处理完再拉取下一条 processMessage(message); } } }).start(); } ``` RocketMQ 也是一样,由消息键 + 组内消费者实例独占队列实现。 ## 消息积压 **消息积压**(Message Backlog)指的是**生产者发送消息的速度超过了消费者处理消息的速度**,导致大量消息在消息队列中堆积未被及时消费 消息积压会带来一系列负面影响: - **业务延迟**:消息处理延迟增加,导致业务操作(如下单、支付通知)响应变慢。 - **系统资源耗尽**:堆积的消息占用大量 Broker 内存和磁盘资源,可能影响 Broker 性能甚至稳定性。 - **用户体验下降**:用户感知到系统响应迟缓,操作卡顿,严重时可能导致服务不可用。 - **数据不一致风险**:如果涉及事务或状态同步,消息积压可能导致数据在一段时间内不一致 消息积压通常由以下一个或多个因素导致: - **消费能力不足**:消费者数量不足、处理逻辑复杂(如涉及慢查询、外部调用)、单条消息处理时间长。 - **生产流量突增**:如促销活动、秒杀场景导致瞬时流量激增,远超消费者处理能力。 - **消费者故障**:消费者应用宕机、网络问题或 Full GC 等导致消费暂停。 - **资源配置不当**:Topic 的队列数不足、消费者线程池配置不合理、QoS/Prefetch 参数设置不当等。 - **Broker 资源瓶颈**:磁盘 I/O、网络带宽或 CPU 资源受限,影响消息读写和分发 ### RabbitMQ 解决 RabbitMQ 解决消息积压的核心思路是**提高消费能力**和**优化资源配置**。思路是:**监控 + 应急方案 + 预防** 1. 首先需要监控队列状态,及时发现积压。 - - **RabbitMQ 管理界面**:查看队列的 `**Ready**`(待消费)、`**Unacked**`(已投递未确认)消息数 - **监控工具**:集成 Prometheus + Grafana 监控队列深度、消费速率等指标,并设置告警(如积压超阈值) 1. 应急处理方案 - - **动态扩容消费者:**临时增加消费者实例数量。通过调整 `**spring.rabbitmq.listener.simple.max-concurrency**` 来动态提高最大并发消费者数(需确保有足够资源) - **启用批量消费****:**调整 `**prefetch**` 值:适当增加 `**prefetch**` 可以让消费者一次接收更多消息(但需权衡处理速度和内存占用) - **临时降级**:如果资源有限,可暂时**暂停非核心业务消费者**或降低其线程数,优先保障核心业务消费 1. 预防消息积压 - - **优化消费者逻辑:**异步化耗时操作;代码中检查慢查询、锁竞争等。 - **合理配置参数:**根据业务场景调整 `**concurrency**`(并发数)和 `**prefetch**`。例如,IO密集型任务可设置较小的 `**prefetch**`(如1-10),计算密集型可设置较大的 `**prefetch**`(如100-200) - **交换机分离**:为不同业务创建**独立的队列和交换机**,避免相互影响。 - **使用死信队列**:对于处理失败的消息,配置死信队列(DLX)进行收集,避免无限重试阻塞队列;定期分析死信队列,修复问题根源。 ### RocketMQ 解决(推荐) RocketMQ 凭借其**高吞吐量**和**灵活的架构**,在处理消息积压方面提供了强大的能力。 思路相同:监控 + 应急方案 + 预防(长期优化方案)。 1. RocketMQ 提供了丰富的监控工具。 - - **RocketMQ 控制台**:直接查看 Topic 的消息堆积量(`**diff**` 值)、消费者组的消费进度(Lag) - **命令行工具**:使用 `**mqadmin**` 命令(如 `**consumerProgress**`、`**topicStatus**`)查询详细状态 - **监控集成**:集成 Prometheus + Grafana,设置积压告警(如 Lag 超过 10 万条) 1. 应急处理方案 - - **动态扩容消费者**:快速启动更多消费者实例。确保**消费者组内的实例数与 Topic 的队列数匹配**(例如,Topic 有 8 个队列,消费者组至少应有 8 个实例)以最大化并行度 - **优化消费线程**:整消费者的线程池大小,增加消费并发度 - **启用批量消费**:启用批量消费模式,减少网络开销。设置 `**consumeMessageBatchMaxSize**`(如每次拉取 32 条消息) 1. 长期优化方案 - - **增加队列数:**如果 Topic 队列数不足成为瓶颈,可通过 `**updateTopic**` 命令增加队列数,提升消费并行度 - **优化 Broker 参数**:将 `**flushDiskType**` 设置为 `**ASYNC_FLUSH**`(异步刷盘)以降低磁盘 I/O 压力 - **优化Broker参数:**调整 `**sendMessageThreadPoolNums**` 和 `**pullMessageThreadPoolNums**` 等线程池参数,提升 Broker 处理能力 - **流量控制与降级**:在紧急情况下,可对生产者进行限流,或对非核心消息启用降级策略(如丢弃部分日志类消息),防止积压进一步恶化 ### 小结 消息积压的解决思路 1. **监控先行**:建立完善的监控和告警机制是发现和解决消息积压的前提。 2. **快速响应**:制定应急预案,包括扩容、降级等流程,以便快速响应积压。 3. **优化消费逻辑**:大多数积压问题与消费逻辑有关,优化消费逻辑是根本。 4. **合理规划**:根据业务峰值进行容量规划,预留缓冲资源。 5. **架构设计**:对于高流量系统,考虑**消息分片**(Sharding)或**多集群**部署,提前分散压力。 ## 事务消息 **事务消息**(Transaction Message)是一种特殊的消息类型,它能确保**生产者发送消息与本地事务的原子性**(要么都成功,要么都失败)。它通过 **两阶段提交(2PC)** 和**事务状态回查**机制,实现了分布式场景下的**最终一致性。** 大白话:生产者给消息队列发送一个半消息(不能被消费的事务消息),生产者在本地把其他事务操作都做了,全部成功,在给消息队列发送一个提交事务的消息,消息队列的半消息就变成普通消息了。 **生命周期:** 半消息 → 本地事务 → 提交/回滚 → 回查确认 → 消费投递。 ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/zrwlDkxPC9KbW4Wx.svg) 事务消息主要为了解决**分布式事务**中**跨服务操作的原子性问题**。在微服务架构中,一个业务操作可能需要跨越多个服务(如订单服务、库存服务、支付服务),可能会发生**数据不一致**,事务消息的作用: - **跨服务操作的一致性**:确保核心业务操作(如订单创建)与多个下游操作(如库存扣减、积分变更、购物车清空)的结果完全一致 - **降低业务侵入**:它将分布式事务的复杂性下沉到消息中间件,业务代码只需关注本地事务和简单的消息确认,开发更简单 举例:假设普通消息执行下面的操作失败了 ```java (1) 创建订单 → 成功 (2) 发送扣库存消息 → 失败 ``` 导致订单系统认为下单成功,但库存系统从没收到扣减请求 → **数据不一致** **如果把这条普通消息设计成事务消息, 让上面的两步变成一个半事务操作:** - **要么订单创建成功 + 消息可靠发送;要么两者都不生效。** 我们来看看RabbitMQ 和 RocketMQ 具体的落地实现 ### RabbitMQ 事务消息 RabbitMQ 的**AMQP协议**提供了原生的事务机制(`**txSelect()**`, `**txCommit()**`, `**txRollback()**`),通过**阻塞式**的方式保证强一致性 **AMQP 协议**: 是一个**二进制的应用层协议**,规定了消息队列系统的通信规范 ,规范中定义了 `tx.select`、`tx.commit`、`tx.rollback` 等命令 ,RabbitMQ 对这三个命令做了具体的实现。事务能力属于AMQP规范的一部分。 **缺点**:性能很差,很少实际使用RabbitMQ 的事务能力,需要事务能力技术选型应该选择 RocketMQ. **特点**: - **强一致性**:通过协议层面保证,可靠性高。 - **阻塞式**:事务提交是同步阻塞的,**性能较低**,吞吐量受限 - **适用场景**:对一致性要求极高的场景,如金融交易。 **工作流程**: 1. **开启事务**:通过 `**channel.txSelect()**` 将信道设置为事务模式。 2. **发送消息**:在事务中发送一条或多条消息。 3. **提交或回滚**:根据业务逻辑执行 `**txCommit()**` 提交事务,或 `**txRollback()**` 回滚事务 ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/ZhXdFWs7f7FY8nud.svg) ### RocketMQ 事务消息(异步最终一致) RocketMQ 事务消息是**高性能、易用**的解决方案,采用**异步化**设计,通过**半消息**和**状态回查**机制实现最终一致性。 **核心概念**: - **半消息**:预先发送的消息,对消费者不可见,状态为“暂不能投递”。 - **本地事务**:生产者执行的业务逻辑(如数据库操作)。 - **状态回查**:Broker 定期回查生产者,确认本地事务的最终状态 **工作流程**: 1. **发送半消息**:生产者发送半消息到 Broker,Broker 持久化后返回 ACK 。 2. **执行本地事务**:生产者执行本地事务(如更新订单状态)。 3. **提交或回滚**:根据本地事务结果,向 Broker 发送 `**COMMIT**` 或 `**ROLLBACK**`。 4. **状态回查**:若 Broker 未收到确认,会定期回查生产者 。 **特点**: - **最终一致性**:不要求强一致性,性能高,适合高并发场景。 - **异步化**:非阻塞,吞吐量高。 - **易用性**:对业务代码的侵入性相对较低。 ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/CXu0PfUKP2k6upkV.svg) 事务消息是分布式事务的一种实现方式。 ## 分布式事务 指一个事务的参与者、涉及的资源服务器分布在不同的物理节点上。它需要保证**跨多个独立资源(数据库、服务)的一系列操作,要么全部成功执行,要么全部失败回滚**,从而维护数据的一致性。 分布式事务主要解决在**分布式系统**中,由于网络分区、节点故障、并发操作等带来的**数据一致性**问题。 实现分布式事务的方案有很多,主流的可以分为两大类:**强一致性方案**和**最终一致性方案**。 ### 强一致性方案(同步阻塞) 这类方案追求所有节点在任何时刻的数据状态都是完全一致的。 - **两阶段提交(2PC - Two-Phase Commit)** - - **角色**:协调者、参与者。 - **第一阶段(准备阶段/投票阶段)**:协调者向所有参与者发送“prepare”请求,询问是否可以执行事务。参与者执行事务操作,**但不提交**,并锁定资源,然后向协调者返回“YES”或“NO”。 - **第二阶段(提交/回滚阶段)**: - - - 如果所有参与者都返回“YES”,协调者发送“commit”请求,所有参与者正式提交事务。 - 如果任何一个参与者返回“NO”或超时,协调者发送“rollback”请求,所有参与者回滚事务。 - - **缺点**:同步阻塞、协调者单点故障、数据锁定时间长、性能差。 两阶段提交是分布式事务的模型,MySQL 的事务是经典实现,三个主流的消息队列没有对应实现。 ------ ### 2. 最终一致性方案(异步非阻塞) 这类方案允许系统在某一时刻数据不一致,但通过补偿机制,最终会达到一致状态。这是目前微服务架构下的主流选择。 #### TCC(Try-Confirm-Cancel) TCC 是一种**分布式事务处理模型**或**设计模式**。它定义了一套标准的业务流程规范,将一个完整的业务逻辑按阶段拆分为三个方法:`**Try**`, `**Confirm**`, `**Cancel**`。任何遵循这套规范的实现,都可以称为 TCC 方案。 实现: TCC 的核心思想是把一个大的分布式事务**拆分成多个独立的本地事务**,每个本地事务都有对应的三个操作。 **核心角色** - **事务发起者**:主业务服务,负责编排整个 TCC 事务流程。 - **事务参与者**:被调用的其他微服务,需要实现 TCC 的三个接口。 - **事务协调器**:通常由框架(如 Seata)承担,负责记录事务全局状态,并驱动 Confirm 或 Cancel 阶段的调用。 **阶段详解** 1. **Try 阶段(资源检查与预留)** - - **目的**:检查业务资源是否可用,并**预留**业务资源。 - **操作**:这个阶段**不执行真正的业务**。例如,在支付场景中,`**Try**` 不是真的扣钱,而是**冻结**账户中的资金;在库存场景中,不是真的减库存,而是**预占**库存。 - **关键点**:Try 操作需要具备**幂等性**,因为可能会被重试。 1. **Confirm 阶段(确认执行业务)** - - **目的**:当所有参与者(多个参与的微服务) `**Try**` 阶段都成功后,执行真正的业务操作。 - **操作**:在支付场景中,`**Confirm**` 才是真正地**扣减**被冻结的资金;在库存场景中,才是真正地**减去**预占的库存。 - **关键点**:Confirm 操作也必须具备**幂等性**。因为网络问题可能导致 Confirm 调用失败,协调器会重试,必须保证多次调用结果与一次调用相同。 1. **Cancel 阶段(取消预留资源)** - - **目的**:如果任何一个参与者(多个参与的微服务)的 `**Try**` 阶段失败,则调用所有已成功 `**Try**` 的参与者的 `**Cancel**` 操作,回滚资源。 - **操作**:在支付场景中,`**Cancel**` 是**解冻**之前冻结的资金;在库存场景中,是**释放**预占的库存。 - **关键点**:Cancel 操作同样必须具备**幂等性**。 **落地实现:** - **Seata**:阿里巴巴开源的分布式事务解决方案,提供了对 TCC 模式的完整支持。它扮演了一个**事务协调者** 的角色,帮助你管理 TCC 事务的生命周期,包括调用 Try、Confirm 或 Cancel 阶段。这是目前国内 Java 生态中最流行的选择。 - **Hmily**:一款高性能的分布式事务 TCC 服务框架,支持多种 RPC 框架,如 Dubbo、Spring Cloud 等。 - **ByteTCC**:一个基于 Java 的 TCC 型事务管理器,兼容 JTA(Java Transaction API)。 ------ #### Saga模式 Saga 模式把一个 **长事务**(跨多个服务/多个数据库)拆分成一系列 **本地事务(Local Transaction)**。 每个本地事务在自己的服务/数据库中执行、提交。 如果某个本地事务失败了,那么之前已经成功的本地事务必须通过 **补偿事务(Compensating Transaction)** 撤销自己的效果。 通常适用于多个服务各自拥有数据库、且需要跨服务的一致性保障,但不能用传统分布式事务(如 2PC)或者不愿承担其性能成本。 落地技术/实现方式包括: 1. **业务逻辑层面定义补偿接口** 在每个本地事务步骤中,同时定义一个“正向操作(Forward Transaction)”和一个“补偿操作(Compensating Transaction)”。例如:库存服务有 `reserveStock()`(扣库存)与 `releaseStock()`(释放库存)作为补偿。 补偿操作通常作为异步任务或消息触发执行 1. **事件驱动或消息队列触发** 在 Saga 模式中,经常使用消息或事件来控制流程:当一个步骤成功后,发布事件触发下一个步骤;当失败时,发布“需要补偿”的事件,触发补偿事务。 1. **Orchestration 或 Choreography 模式控制** - **编排(Orchestration)**: 有一个中央协调器负责记录每一步状态,决定何时执行补偿。 - **编舞(Choreography)**: 各服务基于事件流触发自己的补偿逻辑,无中央协调器。 1. **持久化 Saga 状态与日志** 为了保证补偿事务可靠执行,系统通常会有一个 Saga 日志或状态存储,用来记录每一步是否成功、是否需要补偿。失败时,系统查日志执行补偿。 1. 保证补偿操作幂等性 补偿事务可能会被重复触发,需要保证其幂等性(多次执行效果相同) ------ #### 本地消息表 在一个服务内(本地数据库 + 该服务的消息发送逻辑)把 **“业务操作 + 消息发送”** 放在同一个本地事务中执行。 → 保证:如果业务表更新成功,消息也写入“消息表”;如果业务失败,消息也不会写入。 这样可避免“业务写成功但消息发送失败”导致的 “状态更新” + “通知消息”不一致问题。 不过很容易看出缺陷, 本地消息表的思路是确保每一个服务内的本地事务都是通过的,如果有一个不通过,就不会走到下游 ,但是上游的服务不会回滚, 它的特性是**保证服务内原子**,但**不保证跨服务的原子回滚**。 总结: “本地事务强”、但“跨服务事务弱”。 **具体技术实现原理** 下面是本地消息表模式的大致实现流程: 1. **业务操作 + 消息插入在一个本地事务里** 在服务 A 所使用的数据库内,进行如下操作: ```plain BEGIN TRANSACTION 更新业务表(例如:创建订单) 插入消息表(例如:to_be_sent = true, message_body = {...}) COMMIT ``` 这样保证了:业务数据变更和消息写入是原子的。 1. **后台扫描或变更数据捕获(CDC)发送消息** - - 一个后台任务或专门线程定期扫描 `消息表` 中状态为 `to_be_sent` 的记录。 - 将消息发送到消息队列/中间件。 - 若发送成功,更新消息表为已发送状态(或者删除消息表记录)。 - 若失败,重试或记录失败。 这一机制保证“在业务成功且消息写入后再发送消息”。 1. **消费者消费 + 下一步服务执行** 其他服务订阅队列消息,执行对应操作(例如库存扣减、支付通知等)。 2. **幂等 + 重试机制** - - 消息表发送过程可能重试多次;消费者也可能重复消费,因此系统需设计幂等逻辑。 - 消息可能延迟、重复,需要消费者做好幂等处理。 ### 消费者‑队列 上的负载均衡 **消费者 向 MQ 拉取数据还是MQ向消费者推送数据?** 分情况讨论: | 消息队列 | 消费模型 | 拉/推机制 | 说明 | | ------------ | --------------------------- | ------------------------ | ------------------------------------------------- | | **Kafka** | 消费者组(Pull) | ✅ **主动拉取(Pull)** | 消费者定期向 broker 主动拉取消息 | | **RabbitMQ** | 基于队列(Push) | ✅ **主动推送(Push)** | Broker 主动将消息推送给消费者 | | **RocketMQ** | 消费者组(Push + Pull混合) | ⚙️ **逻辑推送,本质拉取** | Broker 通知消费者有消息,消费者再发 Pull 请求拉取 | **1. RocketMQ 实现** **RocketMQ 的消费者不是被动接收消息,而是 主动拉取(pull) 模式(即使是“推送式”消费者,也是在内部通过定时轮询拉取)。** 负载均衡规则:RocketMQ 支持多种负载均衡规则,默认是**平均分配。** 例如: - Topic:`OrderTopic` - 有 8 个 MessageQueue - Consumer Group:`OrderGroup` - 启动了 4 个消费者实例 8个队列,4个消费者实例,负载均衡分配如下: | Consumer | MessageQueues | | -------- | ------------- | | C1 | MQ0, MQ1 | | C2 | MQ2, MQ3 | | C3 | MQ4, MQ5 | | C4 | MQ6, MQ7 | 这就是**平均分配**,由于队列的消息可能会消费完毕,所以每隔20秒会执行一次重新平均分配。 或者其它触发逻辑: - 消费者实例数量变化(新增或宕机) - Topic 的 MessageQueue 数量变化(比如 Broker 扩容) - 定时任务触发 1. **RabbitMQ 实现** 在 RabbitMQ 中,由队列推送消息给消费者。 多消费者订阅同一个队列时,**默认会轮询(round‑robin)分配消息给各个消费者**。 并且是**带补偿的轮询分发**。 假设两个消费者 C1,C2订阅了队列 order_queue , 那么轮询规则是这样的: | 消息 | 分发给谁 | | ---- | -------- | | msg1 | C1 | | msg2 | C2 | | msg3 | C1 | | msg4 | C2 | | … | … | 消费者消费完毕需要发回一个 ack 给队列,如果队列没有收到 ack , RabbitMQ 会把消息重新分配给别的消费者,确保这条消息被消费然后再消费下一条,这就是带补偿的轮询分发。 #### # 本文参考 https://rocketmq.apache.org/zh/docs/featureBehavior/03fifomessage/ https://rocketmq.apache.org/zh/docs/featureBehavior/04transactionmessage/ https://www.cnblogs.com/wunsiang/p/12765158.html https://juejin.cn/post/7491971657202073654 https://www.cnblogs.com/RunningSnails/p/17373107.html

Cursor 1.2重磅更新,这个痛点终于被解决了!

大家好,我是程序员鱼皮。分享一个重磅消息,AI 编程工具 Cursor 1.2 版本正式发布了! 感觉最近 Cursor 团队像打了鸡血一样,从 [1.0](https://mp.weixin.qq.com/s/3F4EshMIjlcsEdYQnwgrcw) 到 1.1 再到 1.2,短短一个月更新了 2 个正式版本。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/3M4bn79N4RZlT6et.webp) 作为一个深度使用 Cursor 的开发者,我第一时间升级到了 Cursor 1.2 版本,不得不说真是太香了。如果你还没用过 Cursor,那你可能错过了目前最强的 AI 编程工具;如果你已经在用,那这次更新绝对会让你的开发效率再上一个台阶,下面来看看这次都更新了些啥? ![](https://pic.code-nav.cn/post_picture/1601072287388278786/qpB0MzBwZrnh0iBB.webp) ## Agent 规划能力 AI Agent 如果要完成复杂的任务,通常会先思考规划如何完成任务、然后再一步步执行。但是之前 AI 的计划对我们来说不够透明,比如让它生成一个复杂的网站,可能除了查看它的思考过程外,你并不知道 AI 总共要做哪些事情、接下来要做什么、当前执行到第几步。 这次更新支持了 `Agent To-dos`,当你给 Agent 一个复杂任务时,它会自动分解成多个子任务,并且清楚地展示任务之间的依赖关系。 比如我让 AI 帮忙写一篇 10w 字的长篇小说,可以看到 AI 生成了有 11 条任务的 To-dos 列表: ![](https://pic.code-nav.cn/post_picture/1601072287388278786/jTsPq0rHGejDOqZV.webp) 是不是清晰很多,一下子就 get 到了接下来 AI 要干什么,能够让你更好地控制任务执行过程。比如我对它的规划不满意,就让它重新规划: ![](https://pic.code-nav.cn/post_picture/1601072287388278786/UqQSsqUVT93yqMvs.webp) 注意,想使用这个功能,需要确保设置中开启了 To-Do List: ![](https://pic.code-nav.cn/post_picture/1601072287388278786/Ph5H6UDCNoBs25Nn.webp) 而且经过我的测试,目前不是所有的提示词都会触发 To-Do List,比如我让它生成一个复杂的网站项目,它就不会规划出任务列表。但如果在提示词中添加 “先规划任务”,就更容易触发。 ## 消息队列 以前使用 Agent 最痛苦的就是等待。你想到一个新需求,但 Agent 还在处理上一个任务,只能干等着。现在有了 **消息队列功能**,你可以直接把后续的指令发送给 Agent,它会自动排队执行。 这个功能对我这种思维跳跃的程序员来说简直太实用了,举个例子,我想修复网站的 10 个 Bug、并且给网站加 5 个新功能。以前我需要一个个提交任务、然后每隔 1 分钟左右再来检查下任务完成情况,再输入下一个任务,很浪费时间,我还没办法中途分心去做别的事。之前我的解决方案是多开几个 Cursor 窗口、或者单独开一个文档记录自己接下来要执行的提示词。而现在有了消息队列,我一次提交十几个任务,然后就可以安心摸鱼去了,过个二十分钟再来整体验收。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/uBngEi5LaN903Lrj.webp) 注意,想使用这个功能,需要确保设置中开启了 Queue Messages: ![](https://pic.code-nav.cn/post_picture/1601072287388278786/2XUitb3KNQQCW1Eu.webp) ## 记忆功能正式上线 Cursor 1.0 的时候推出了 Memories 功能,这次它终于转正了。 这个记忆功能和上下文对话历史(也就是聊天记录)是有区别的,不是什么都记,更多的是 **记忆规则**。比如你经常使用某种代码风格,或者有特定的生成项目的要求,Cursor 会自动记住这些信息,在后续生成时主动应用。 举个例子,我这里让 AI 以后尽量用 Windows 系统的命令来生成代码。执行后,可以看到记忆被更新了,里面是一条规则。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/2hJEJvtWKvQ3uJG5.webp) 之后在这个项目中生成代码时,就会使用这个规则。还可以在规则设置页面进行管理和删除。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/A8JSpLvGkHvVYKuX.webp) 这样一来,通过持续不断地对话,AI 助手会变得越来越了解你。 ## PR 索引和搜索 新增的 PR 索引和搜索功能可以让代码审查变得更加智能。Cursor 现在可以: - 自动索引和总结 Pull Request - 语义化搜索历史 PR - 关联 GitHub 评论和 BugBot 审查结果 - 支持 Slack 集成,方便团队协作 当你需要排查某个 Bug 时,AI 可以直接搜索相关的历史 PR,快速定位问题根源,对于维护大型项目应该会挺有帮助的。 注意,想使用这个功能,需要确保设置中开启了对 PR 的索引: ![](https://pic.code-nav.cn/post_picture/1601072287388278786/CGsv53CSaY7wMc3b.webp) 不过我试了很多次,都没有触发官方演示的那种 PR 读取效果,反而 AI 会利用 git 命令来查找提交记录,看来意图识别准确度还要再继续优化优化。 ![](https://pic.code-nav.cn/post_picture/1601072287388278786/ln8D6HwqhX89w39B.webp) ## 更快的代码补全 Tab 补全速度提升了约 100ms,首次响应时间减少了 30%。别小看这 100ms,对于高频使用代码补全功能的朋友来说,这个优化能够明显提升编程体验的流畅度。 官方提供的性能对比图: ![](https://pic.code-nav.cn/post_picture/1601072287388278786/kIuWkwCy0e4c2rIR.webp) ## 智能冲突解决 当出现代码合并冲突时,Agent 现在可以 **自动尝试解决冲突**。点击 “在聊天中解决”,相关上下文会被自动添加到对话中,Agent 会分析冲突原因并提供解决方案。 ------ 还有一些其他改动,比如代码库搜索使用了新的嵌入模型来提高准确度;还有 Background Agent 的一些优化。这些也不需要我们关心。 ## 总结 总的来说,这次更新对我来说最有用的功能是消息队列。我相信很多朋友也和我一样,随着 AI 的发展,越来越依赖 AI,工作内容从独立思考 + 执行变成了等着 AI 返回内容,等待的过程中也不知道自己在干嘛,不知不觉时间就过去了。这个功能真的解决了我经常要等待 AI、被 AI 打断工作的痛点,也期待 AI 编程工具接下来都能朝着更加智能、更加人性化的方向发展,让 Vibe Coding 流行起来! 大家有没有用过 Cursor?对这次更新有什么看法?欢迎在评论区分享,对 AI 感兴趣的朋友可以免费获取 [鱼皮开源的 AI 知识库](https://github.com/liyupi/ai-guide)。 ## 更多 💻 编程学习交流:[编程导航](https://www.codefather.cn/) 📃 简历快速制作:[老鱼简历](https://laoyujianli.com) ✏️ 面试刷题神器:[面试鸭](https://mianshiya.com) 📖 AI 学习指南:[AI 知识库](https://ai.codefather.cn/)

深入浅出 RocketMQ 消息队列 笔记(13)

如何处理消息堆积问题 消息堆积的本质:消费速度<生产速度 常见消息堆积原因 1.瞬时流量 2.消费性能不足 3.机器不够 4.Bug 5.其他功能影响 瞬时流量 比如平常一分钟一条消息,瞬时流量来了,可能一秒钟就产生了上万条消息。 这里要对消息的瞬时性做分析,如果只是瞬时的,持续几分钟就好了,那就不需要做任何改动,毕竟消息队列就是负责削峰填谷的。如果持续时间长且消息量级过大,那堆积造成的影响还是需要重视。 这里举个例子,是yes哥之前遇到的一个情况,有七千万的历史数据需要清洗(消费),其中生产者花了12天时间将数据发给消息队列,但消费者一天只能消费200w条消息,所以乐观估计,消费者消费完所有历史数据需要至少35天。 这里会带来不少问题,如果消息队列用的是第三方提供的(比如阿里云),消息存储有个默认时间3天,如果消息存储3天还没被消费,那这个消息就会被删掉了(删掉的目的是为了控制存储成本),堆积未消费的消息丢失后,还需要生产者重新发送一遍,会很麻烦 还有一个问题,如果历史数据跟正常业务数据的处理流程是相同的,也就是他们的topic相同,那在broker产生堆积后,光顾着处理历史数据消息了,正常业务数据的处理也会因此而堵住了。 ![image.png](https://pic.code-nav.cn/post_picture/1852176739385810946/JLk0cyiBGpGlABZn.webp) 这里yes哥的解决方式是降低生产消息的频率,不持续发消息,让消息发送的频率略低于消费速率,给正常业务数据有消费的空间 消费性能问题 1.可以把循环里的查询拎出来放在外面批量查询,然后Map映射 2.单条插入/更新改成批量插入/更新 机器不够 消费性能上能优化还是优化了,但消息还是堆积,那就得水平扩展机器了,不过要注意,消费者与队列数要一直保持消费<队列的数量,否则会造成消费者空忙,不能很好的利用重平衡机制。 ![image.png](https://pic.code-nav.cn/post_picture/1852176739385810946/XXYS26WKmSulTu3v.webp) Bug 也有可能是消费逻辑有bug,导致消息堆积。消费失败需要重试,默认重试16次才会进入死信队列,相当于一条消息用到的流量被放大了16倍,资源都用来重试了,重试后还是失败。再如果消息是顺序消息,那后面的消息都处理不下去了。这种情况只能对消息队列做监控,发现报错后紧急发布版本,修复bug。进入死信队列的消息可以手动重新消费。 其他功能影响 其他业务功能抢夺了消息队列在使用的资源,比如消息消费的逻辑是扣减库存,其他业务功能也涉及到扣减库存,那么就有可能会争夺同一个商品的处理,也就是同一个数据库的行锁,导致消费速率降低。还有可能是其他功能的长事务导致消费等待,等等。如果平常正常,流量也没有变很大,有天突然堆积了,可以在这方面考虑问题原因。

深入浅出 RocketMQ 消息队列 笔记(12)

如何保证消息不重复 消息无法保证不重复,但可以保证消息被幂等消费,等同于仅消费一次 消息不重复解决方案 1.项目初期/新功能开始设计的时候,就考虑好幂等设计,满足消息的幂等消费 2.调整业务执行顺序,将影响大业务幂等的逻辑提前,提前终止重复消费 3.添加唯一性索引,如果有业务上的唯一索引(比如订单号)就不用加,反之就加一个唯一索引,重复消费插入时会报错的。如果不方便添加唯一索引,可以再开一张表存唯一索引,然后将这两张表放在同一个事务里 4.引入第三方中间件,比如redis,处理的时候用SETNX判断,已插入就直接返回,否则执行正常业务逻辑。不过这里有可能在redis刚存值的时候系统宕机了,导致消息被假处理,这个需要注意。

深入浅出 RocketMQ 消息队列 笔记(11)

保证消息不丢失的必要条件 生产者发送消息、生产者存储消息、消费者拉取消息,需要保证三大流程消息不丢失,缺一不可 生产者保证消息完整发送并存储至broker broker保证存储的消息不丢失 消费者保证拉取的消息一定被消费,即使重启了,也能确保未消费的消息继续消费 生产+发送消息流程 以下单为例,下单后增加积分,增加积分这个动作放在消息里实现,且要保证该消息一定发送成功。 这里可以引用TCP协议的ack请求确认机制。如果broker收到生产者推送来的消息,就返回ack给生产者,这样就能保证生产者发送消息这个阶段,消息不会丢失。如果生产者一直没收到ack,可能是网络或者其他原因,这时候就会进行重试,重试次数达到上限后抛异常。但这里抛异常不能影响主流程,所以需要对异常进行特殊处理。如果是同步发送消息,那就try-catch,如果是异步,那就在异步方法对应的onException做异常处理,这里可能要做一些补偿机制。 存储流程 broker在返回ack给生产者之前要确保消息已经成功存储了,RocketMQ的消息默认是异步刷盘,先刷到cache上,再等操作系统或定时刷盘任务把消息刷到磁盘上。如果此时断电了,消息就丢失了。所以为了保证消息不被丢失,这里可以把刷盘方式改成同步刷盘(flushDiskType=SYNC_FLUSH)。如果broker是集群,那也得保证主从broker的复制方式是同步复制,这样的话消息就更安全了。 消费流程 消费者消费后需要上报点位给broker,告诉broker已经消费到第几条消息了。如果消费者在处理消息的时候异步处理,然后直接告知broker消费成功,就很有可能在消费过程中报错,导致再次拉取消息的时候,从之前上报过的点位继续拉消息,消息就“丢失”了。所以这里要确保消费者真正消费完成消息后,再提交点位。

深入浅出 RocketMQ 消息队列 笔记(10)

RocketMQ的消息存储 用一个CommitLog所有分发给该broker的消息,多Topic混合存储 commitLog超过1G,会新起一个commitLog 每条消息存储到commitLog都会在consumeQueue生成一条记录,可以视为一个索引(稠密索引) Kafka的消息存储 Kafka在Topic下也分了多个队列来提高消费的并发度,但在Kafka不叫队列,叫分区partition Kafka的消息存储和RocketMQ略有不同,Kafka的消息存储以Partition为单位进行存储 每个Topic的每个分区都有自己的消息文件、索引文件和时间索引文件,他们的文件名相同,后缀名不同 文件名的命名规则是第一条消息 的offset,文件写满会新起一个文件 索引文件设计的和RocketMQ不同,Kafka是每隔几条消息再创建一条索引,节省了存储空间,能保存更多的索引,这样的索引叫稀疏索引 稀疏索引如何找到对应消息 通过offset找到对应的索引文件,通过二分遍历找到离消息最近的索引,再通过这个索引找到消息文件里此条消息的位置,再遍历消息文件找到目标消息 Kafka时间复杂度:O(log2n)+O(m),n为索引个数,m为稀疏程度 RocketMQ时间复杂度:O(1) 所以这里就需要权衡利弊了,Kafka是时间换空间,RocketMQ是空间换时间 虽然Kafka是按分区存储消息,但这样可能会引起 Kafka的存储设计对数据复制和迁移很友好,但在海量Topic、Partition场景会有性能问题

深入浅出 RocketMQ 消息队列 笔记(9)

Broker集群 单master 只有一个broker,挂了的话消息队列就不能用了 多master 一个master挂了的话,生产者会往其他正常的master发消息,但原先挂了的broker存储的消息得等重启后才能继续消费了。这里存储还涉及到异步刷盘和同步刷盘。 异步刷盘 性能高一点,但有可能会消息丢失 同步刷盘 性能低一点,但消息不会丢 不过多master只能对普通消息有用,顺序消息的话,顺序性就无法保证了 多master多slave异步复制 slave一般是拿来backup,不会拿来主动接收消息 多master多slave同步复制 消息可靠性比较好,但性能也有影响 Dledger 指一组具有相同名称的Broker,至少有3个节点,组成RocketMQ-on-Dledger Group(基于一致性协议raft),当主节点挂了,会自动选举出新的主节点 可以理解为解决了顺序消息的问题,集群中的消息存储都是一样的,不会出现消息顺序不一致等数据冲突问题,leader挂了,follower直接替上就好了

MQ总结 v1.0

1. 什么时候用? --------- 场景: 1. 系统处理高耗时 需要异步操作 2. 解决原来的线程池丢消息的情况 2. 使用 ----- ### 2.1. 开始 1. 安装mq 2. 安装mq管理UI 初始账户密码为guest 3. 导入依赖 编写demo 测试 1. 创建工厂 建立连接 2. 创建信道 用信道声明队列 声明交换机 绑定交换机等操作 **生产者消费者在声明队列时需参数保持一致** ### 2.2. 交换机 exchange 介绍: 如计网里面的交换机根据MAC地址的转发数据帧的效果类似 不同的交换机可以根据不同的规则转发消息 分类: fanout; direct; topic; #### 2.2.1. fanout 介绍: '**扇出' 交换机 名如其机 感觉一下子消息全部扇到了脸上一样** **特点: 就像计网的广播帧 转发给所有信道 每个都可以收到消息** <img src="https://pic.code-nav.cn/post_picture/1608648175386624001/NzeNmC83YC1YyAlo.png" alt="" width="433px" /> 使用:[https://www.rabbitmq.com/tutorials/tutorial-three-java](https://www.rabbitmq.com/tutorials/tutorial-three-java) 生产者信道在声明交换机时 就声明fanout 所有绑定该交换机的消费者信道都会受到消息 #### 2.2.2. direct 介绍: 直连交换机 可以根据路由key 直接发送消息 单绑定: ![](https://pic.code-nav.cn/post_picture/1608648175386624001/GkTUegVS7t0eFLRb.png) 多绑定: ![](https://pic.code-nav.cn/post_picture/1608648175386624001/1j0YVoFmDpSCtZsi.png) 使用:[https://www.rabbitmq.com/tutorials/tutorial-four-java](https://www.rabbitmq.com/tutorials/tutorial-four-java) 在声明队列时 声明路由参数 那么创建的交换机在转发消息时就会根据路由key发送到相应的队列 可以应用的场景: **根据不同的路由 发送不同的日志消息到消费者** ![](https://pic.code-nav.cn/post_picture/1608648175386624001/cbwrhZ2A3nCqnM1f.webp) #### 2.2.3. topic 介紹: topic 交换机是用根据规则来模糊匹配路由key的交换机 ![](https://pic.code-nav.cn/post_picture/1608648175386624001/oFQYsdi2nQvWTKqu.png) <img src="https://pic.code-nav.cn/post_picture/1608648175386624001/6Ybx07r1yq7U2ImB.png" alt="" width="625px" /> *是用来匹配一个词 #是用来匹配0个或者多个词 **举例:** ![](https://pic.code-nav.cn/post_picture/1608648175386624001/jgrLSLWrf1M7n16w.png)Q1: *.orange.* Q2: *.*.rabbith 和`lazy.#`. 那么 一下方式 quick.orange.rabbit 匹配Q1 Q2 lazy.orange.elephant 匹配Q1 Q2 quick.orange.fox 匹配Q1 lazy.brown.fox 匹配Q2 lazy.pink.rabbit 匹配Q2 只匹配一次 使用 [https://www.rabbitmq.com/tutorials/tutorial-five-java](https://www.rabbitmq.com/tutorials/tutorial-five-java) ### 2.3. 发布确认机制 Publish Confirm #### 2.3.1. 设置消息确认的到期时间 ```java while (thereAreMessagesToPublish()) { byte[] body = ...; BasicProperties properties = ...; channel.basicPublish(exchange, queue, properties, body); // uses a 5 second timeout channel.waitForConfirmsOrDie(5_000); } ``` 在规定的到期时间内 未确认消息 或者该消息被拒绝 就会抛出一个异常 处理该异常通常是 重发或者是打印错误日志 **优点: 简单** **缺点: 当确认消息阻塞发布消息时 发布效率就会很低 不适合高吞吐量的系统** #### 2.3.2. 批量发布 ```java // 设置批次数量大小 int batchSize = 100; int outstandingMessageCount = 0; while (thereAreMessagesToPublish()) { byte[] body = ...; BasicProperties properties = ...; channel.basicPublish(exchange, queue, properties, body); outstandingMessageCount++; // 发送数量没到就不设置过期时间 if (outstandingMessageCount == batchSize) { channel.waitForConfirmsOrDie(5_000); outstandingMessageCount = 0; } } // 发送完了 发现都还没到数量 就直接设置过期时间 if (outstandingMessageCount > 0) { channel.waitForConfirmsOrDie(5_000); } ``` 优点: 相比于接受单独的确认 效率好很多 缺点: 出问题不知道是哪个出了错 #### 2.3.3. 手动确认 ### 2.4. 死信 死信: 1. 被拒绝的消息 2. 过期的消息 3. 被删除的消息 当定义了死信交换机和死信队列 这些消息就会由该交换机转发到死信队列 #### 2.4.1. 死信交换机 设置: [https://www.rabbitmq.com/docs/dlx](https://www.rabbitmq.com/docs/dlx) ```java channel.exchangeDeclare("some.exchange.name", "direct"); // Important: prefer using policies over hardcoded x-arguments Map<String, Object> args = new HashMap<String, Object>(); //绑定个死信交换机 args.put("x-dead-letter-exchange", "some.exchange.name"); //给绑定的死信交换机 指明它所要转发死信到哪条死信队列上去 args.put("x-dead-letter-routing-key", "some-routing-key"); //此时 该队列就绑定了死信交换机 channel.queueDeclare("myqueue", false, false, false, args); ``` 死信交换机就像direct交换机一样 正常创建就好 #### 2.4.2. 示例 ![](https://pic.code-nav.cn/post_picture/1608648175386624001/sdjGWd3odiaOf9IZ.webp) 创建死信的direct交换机和队列 ```java public class DieLetterWorker { private final static String QUEUE_NAME1 = "m3"; private final static String QUEUE_NAME2 = "m4"; public static void main(String[] argv) throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("localhost"); Connection connection = factory.newConnection(); Channel channel1 = connection.createChannel(); Channel channel2 = connection.createChannel(); //声明死信交换机 channel1.exchangeDeclare("e2", "direct"); //声明死信队列 channel1.queueDeclare(QUEUE_NAME1, true, false, false, null); //绑定队列到交换机 channel1.queueBind(QUEUE_NAME1, "e2", "worker1"); //声明交换机 同上 channel2.exchangeDeclare("e2", "direct"); channel2.queueDeclare(QUEUE_NAME2, true, false, false, null); channel2.queueBind(QUEUE_NAME2, "e2", "worker2"); System.out.println(" [*] Waiting for messages. To exit press CTRL+C"); DeliverCallback deliverCallback1 = (consumerTag, delivery) -> { try { String message = new String(delivery.getBody(), StandardCharsets.UTF_8); System.out.println(" [x]worker1 Received1 '" + message + "'"); //模拟消息处理 Thread.sleep(1000); //如果消息处理失败,将消息发送到死信队列 } catch (InterruptedException e) { throw new RuntimeException(e); } }; DeliverCallback deliverCallback2 = (consumerTag, delivery) -> { try { String message = new String(delivery.getBody(), StandardCharsets.UTF_8); System.out.println(" [x]worker2 Received1 '" + message + "'"); //模拟消息处理 Thread.sleep(1000); } catch (InterruptedException e) { throw new RuntimeException(e); } }; channel1.basicConsume(QUEUE_NAME1, true, deliverCallback1, consumerTag -> { }); channel2.basicConsume(QUEUE_NAME2, true, deliverCallback2, consumerTag -> { }); } } ``` 生产者 ```java public class DieLetterProvider { private final static String QUEUE_NAME1 = "hello1"; public static void main(String[] argv) throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("localhost"); try (Connection connection = factory.newConnection(); //创建两个消息队列 Channel channel = connection.createChannel(); ) { //声明交换机 channel.exchangeDeclare("e1", "direct"); //手动输入 Scanner scanner = new Scanner(System.in); while (scanner.hasNextLine()) { String message = scanner.nextLine(); String[] strs = message.split(" "); //将消息发送到交换机 不指定路由 channel.basicPublish("e1", strs[1], null, message.getBytes(StandardCharsets.UTF_8)); System.out.println(" [x] Sent '" + message + "'"); } } } } ``` 消费者 ```java public class DieLetterConsumer { private final static String QUEUE_NAME1 = "m1"; private final static String QUEUE_NAME2 = "m2"; public static void main(String[] argv) throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("localhost"); Connection connection = factory.newConnection(); Channel channel1 = connection.createChannel(); Channel channel2 = connection.createChannel(); //声明死信交换机 channel1.exchangeDeclare("e2", "direct"); channel2.exchangeDeclare("e2", "direct"); //声明死信队列参数 HashMap<String, Object> args1 = new HashMap<>(); args1.put("x-dead-letter-exchange", "e2"); args1.put("x-dead-letter-routing-key", "worker1"); HashMap<String, Object> args2 = new HashMap<>(); args2.put("x-dead-letter-exchange", "e2"); args2.put("x-dead-letter-routing-key", "worker2"); //声明交换机 channel1.exchangeDeclare("e1", "direct"); //声明队列 channel1.queueDeclare(QUEUE_NAME1, true, false, false, args1); //绑定队列到交换机 channel1.queueBind(QUEUE_NAME1, "e1", "exc1"); //声明交换机 同上 channel2.exchangeDeclare("e1", "direct"); channel2.queueDeclare(QUEUE_NAME2, true, false, false, args2); channel2.queueBind(QUEUE_NAME2, "e1", "exc2"); System.out.println(" [*] Waiting for messages. To exit press CTRL+C"); DeliverCallback deliverCallback1 = (consumerTag, delivery) -> { try { String message = new String(delivery.getBody(), StandardCharsets.UTF_8); if (message.equals("hello exc1")) { throw new InterruptedException("hello exc1"); } System.out.println(" [x]exc1 Received1 '" + message + "'"); //模拟消息处理 Thread.sleep(1000); //如果消息处理失败,将消息发送到死信队列 //确认消息 channel1.basicAck(delivery.getEnvelope().getDeliveryTag(), false); } catch (InterruptedException e) { //如果消息处理失败,将消息发送到死信队列 channel1.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false); } }; DeliverCallback deliverCallback2 = (consumerTag, delivery) -> { try { String message = new String(delivery.getBody(), StandardCharsets.UTF_8); if (message.equals("hello exc2")) { throw new InterruptedException("hello exc2"); } System.out.println(" [x]exc2 Received1 '" + message + "'"); //模拟消息处理 Thread.sleep(1000); //确认消息 channel2.basicAck(delivery.getEnvelope().getDeliveryTag(), false); } catch (InterruptedException e) { //如果消息处理失败,将消息发送到死信队列 channel2.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false); } }; channel1.basicConsume(QUEUE_NAME1, false, deliverCallback1, consumerTag -> { }); channel2.basicConsume(QUEUE_NAME2, false, deliverCallback2, consumerTag -> { }); } } ```

下载 APP