ROCKETMQ-消息发送与消费(二)

💡 原文中文,约9400字,阅读约需23分钟。
📝

内容提要

本文介绍了ROCKETMQ的消息发送和消费相关关注点,包括可靠发送消息、将消息发送到broker、消息发送类型和行列选择,以及中心类和推拉形式的使用示例。还讨论了消息消费的一致性和并发消费时的问题。

🔎

延伸解读

消息发送的故障规避机制

文章指出,非顺序消息发送默认采用轮询策略,并引入重试机制(默认重试2次)。但单纯轮询在Broker故障时可能连续失败,因此加入故障规避机制:重试时尽量避开上次发送失败的Broker,转而选择其他Broker上的队列,以提高发送成功率。该机制由sendLatencyFaultEnable参数控制,是提升消息发送高可用性的重要设计。

推模式与拉模式的取舍

RocketMQ提供推(Push)和拉(Pull)两种消费模式。推模式(DefaultMQPushConsumer)封装了拉模式的复杂性,简化了用户使用,适合大多数场景。拉模式(DefaultMQPullConsumer)API较底层,使用不便,尤其在多消费者场景下更为复杂。为此,4.6.0版本引入DefaultLitePullConsumer,提供更简便的API,并支持通过pullThreadNums参数设置拉取线程数(默认20),这是Lite Pull相比Push模式的一个显著优势。

顺序消息的消费保证与风险

顺序消息要求发送到同一队列并按顺序消费。消费端通过加锁保证单队列的顺序处理,且在Broker端请求队列锁,确保同一时刻只有一个消费者能拉取该队列。然而,并发消费的重试机制(默认重试16次后进入死信队列)会破坏顺序性,因此顺序消费的重试次数为Integer.MAX_VALUE,在消费端不断重试。若一条消息始终无法消费成功,消费进展将无法推进,导致消息积压。

消费进展提交与重复消费

在集群模式下,消费进展存储在Broker端;广播模式则存储在消费者本地。提交消费进展时,RocketMQ取ProcessQueue中最小的偏移量作为进展。这可能导致重复消费:例如消息1~5中,3、4和1已消费,但2、5未完成,提交偏移量为2。若此时客户端挂掉,重启后会从消息2开始消费,导致3、4被重复消费。这是使用推模式时需要注意的消费语义问题。

❓

Q&A

ROCKETMQ支持哪些消息发送类型?

ROCKETMQ支持同步发送、异步发送、单向发送和批量发送。

如何提高ROCKETMQ消息发送的高可用性?

通过引入重试机制和故障规避机制来提高消息发送的高可用性。

ROCKETMQ的消息消费模式有哪些?

ROCKETMQ的消息消费模式包括推模式和拉模式。

什么是DefaultMQPushConsumer?

DefaultMQPushConsumer是ROCKETMQ中推模式的默认实现类,简化了用户的使用。

ROCKETMQ如何处理事务消息?

ROCKETMQ通过保证与数据库事务的一致性来处理事务消息,可以使用本地事务表方法。

顺序消息在ROCKETMQ中如何保证消费顺序?

顺序消息需要保证在同一队列中按顺序消费,消费端需加锁以确保顺序性。

🏷️

标签

➡️

继续阅读