MQ

各种MQ的比较

MQ有rabbitMq,rocketMq,kafka,activeMq

image.png
image.png
image.png
image.png

MQ有什么优缺点

导致系统复杂性增加

rocketMq 的组成

image.png

生产者(Producer):负责产生消息,生产者向消息服务器发送由业务应用程序系统生成的消息。
消费者(Consumer):负责消费消息,消费者从消息服务器拉取信息并将其输入用户应用程序。
消息服务器(Broker):是消息存储中心,主要作用是接收来自 Producer 的消息并存储, Consumer 从这里取得消息。
名称服务器(NameServer):用来保存 Broker 相关 Topic 等元信息并给 Producer ,提供 Consumer 查找 Broker 信息

MQ的整体流程

image.png
image.png

http://svip.iocoder.cn/RocketMQ/Interview/#%E8%AF%B7%E8%AF%B4%E8%AF%B4%E4%BD%A0%E5%AF%B9-Consumer-%E7%9A%84%E4%BA%86%E8%A7%A3%EF%BC%9F
https://blog.csdn.net/javahongxi/article/details/84931747

MQ怎么保证消息不丢失

Producer发送消息的时候回接受ack 回馈
最大重试次数3次
Customer消费消息的幂等性以及消费消息的处理

https://zhuanlan.zhihu.com/p/161965554

手段一:提供sync的发消息的方式,等待broker处理结果,提供了三种方式
1、 同步发送,阻塞当前线程等待broker的响应发送结果
2、异步发送,producer先创建一个发送给broker消息的任务,把该任务提交给线程池,等执行完该任务时,回调用户自定义的回调函数,执行处理结果。
3、 Oneway发送,oneway只负责发送请求,不等待应答,Producer只负责把请求发出去,而不处理响应结果。
我们在调用producer.send方法时,不指定回调方法,则默认采用同步发送消息的方式,这也是丢失几率最小的一种发送方式。
手段二:发送消息如果失败或者超时,则重新发送。
•发送重试源码如下,本质其实就是一个for循环,当发送消息发生异常的时候重新循环发送。默认重试3次,重试次数可以通过producer指定。
手段三:broker提供多master模式,即使某台broker宕机了,保证消息可以投递到另外一台正常的broker上。
• 如果broker只有一个节点,则broker宕机了,即使producer有重试机制,也没用,因此利用多主模式,当某台broker宕机了,换一台broker进行投递。

总结
•producer消息发送方式虽然有3种,但为了减小丢失消息的可能性尽量采用同步的发送方式,同步等待发送结果,利用同步发送+重试机制+多个master节点,尽可能减小消息丢失的可能性。

Broker处理消息阶段

public enum FlushDiskType { SYNC_FLUSH, //同步刷盘 ASYNC_FLUSH//异步刷盘(默认) }
我们知道,当消息投递到broker之后,会先存到page cache,然后根据broker设置的刷盘策略是否立即刷盘,也就是如果刷盘策略为异步,broker并不会等待消息落盘就会返回producer成功,也就是说当broker所在的服务器突然宕机,则会丢失部分页的消息。
手段五:提供主从模式,同时主从支持同步双写
1、即使broker设置了同步刷盘,如果主broker磁盘损坏,也是会导致消息丢失。 因此可以给broker指定slave,同时设置master为SYNC_MASTER,然后将slave设置为同步刷盘策略。
此模式下,producer每发送一条消息,都会等消息投递到master和slave都落盘成功了,broker才会当作消息投递成功,保证休息不丢失。

RocketMQ默认broker的刷盘策略为异步刷盘,如果有主从,同步策略也默认的是异步同步,这样子可以提高broker处理消息的效率,但是会有丢失的可能性。因此可以通过同步刷盘策略+同步slave策略+主从的方式解决丢失消息的可能。

Consumer消费消息阶段

从producer投递消息到broker,即使前面这些过程保证了消息正常持久化,但如果consumer消费消息没有消费到也不能理解为消息绝对的可靠。因此RockerMQ默认提供了At least Once机制保证消息可靠消费。
何为At least Once?
Consumer先pull 消息到本地,消费完成后,才向服务器返回ack。
通常消费消息的ack机制一般分为两种思路:
1、先提交后消费;
2、先消费,消费成功后再提交;
思路一可以解决重复消费的问题但是会丢失消息,因此Rocketmq默认实现的是思路二,由各自consumer业务方保证幂等来解决重复消费问题。
手段七:消费消息重试机制
当消费消息失败了,如果不提供重试消息的能力,则也不能算完全的可靠消费,因此RocketMQ本身提供了重新消费消息的能力。
总结
consumer端要保证消费消息的可靠性,主要通过At least Once+消费重试机制保证。

MQ的高可用

image.png
image.png
image.png
image.png
image.png
image.png

https://zhuanlan.zhihu.com/p/161965554
https://zhuanlan.zhihu.com/p/337530036(知乎MQ的入门)

消息消费的顺序性

image.png

image.png

🦅 实现原理
顺序消息的实现,相对比较复杂,想要深入理解的胖友,可以看看 《RocketMQ 源码分析 —— Message 顺序发送与消费》 。
具体的代码实现,可以看看 《芋道 Spring Boot 消息队列 RocketMQ 入门》的「8. 顺序消息」 小节。

MQ 的积压

一个消费过慢,一个是生产过快,如果生产过快,可以进行一些限流
如果消费过慢,可能是程序出问题,手动的去解决,或者机器过少,增加机器等等

©著作权归作者所有,转载或内容合作请联系作者
【社区内容提示】社区部分内容疑似由AI辅助生成,浏览时请结合常识与多方信息审慎甄别。
平台声明:文章内容(如有图片或视频亦包括在内)由作者上传并发布,文章内容仅代表作者本人观点,简书系信息发布平台,仅提供信息存储服务。

相关阅读更多精彩内容

  • MQ(Message Queue)是一种跨进程的通信机制,用于传递消息。是一种高效、可靠、安全、可扩展的分布式消息...
    ShawnCaffeine阅读 1,151评论 0 0
  • 一. 认识消息队列 1. 队列 队列(queue)是只允许在一端进行插入操作,而在另一端进行删除操作的线性表(数据...
    Serializable_dx阅读 2,340评论 0 1
  • 为什么使用消息队列 解耦 看这么个场景。A 系统发送数据到 BCD 三个系统,通过接口调用发送。如果 E 系统也要...
    程序猿TODO阅读 468评论 1 0
  • MQ消息中间件,面试能问些什么? 为什么使用消息队列?消息队列的优点和缺点? kafka、activemq、rab...
    码农开花阅读 389评论 0 1
  • 个人专题目录[https://www.jianshu.com/u/2a55010e3a04] 5. 消息中间件MQ...
    Java及SpringBoot阅读 633评论 0 0

友情链接更多精彩内容