跳转至主要内容
版本: 5.0

消费进度管理

Apache RocketMQ 使用消费者位点(Consumer Offset)来管理消费者的进度。本主题介绍了 Apache RocketMQ 的消费者进度管理机制。

背景

在 Apache RocketMQ 中,消息可以在消费者订阅之前或之后产生。那么消费者如何知道从哪里开始消费消息,以及已消费的消息是如何标记的呢?为了解决这一挑战,Apache RocketMQ 开发了消费者进度管理机制。

Apache RocketMQ 的消费者进度管理机制解决了以下问题:

  • 客户端启动后从哪里开始消费消息?

  • 如何标记已消费的消息以确保其不会被多次处理?

  • 如果发生服务异常,同一客户端能否再次消费该消息?

工作原理

消息位点 (Message Offset)

在 Apache RocketMQ 中,消息按照到达顺序排队存放在 Topic 中,并被分配一个唯一的 Long 型坐标,这被称为消息的位点(Offset)。有关这些概念定义的更多信息,请参阅 Topic消息队列

理论上,一个消息队列可以存储无限数量的消息。因此,Offset 的取值范围是从 0 到 Long.MAX_VALUE。您可以根据 Topic、队列和 Offset 定位任何消息。下图显示了这三个概念之间的关系。Offset

在 Apache RocketMQ 中,队列中最早消息的位点称为最小位点(MinOffset),最新消息的位点称为最大位点(MaxOffset)。尽管消息队列理论上可以容纳无限数量的消息,但存储它们的物理机空间是有限的。因此,Apache RocketMQ 会动态删除队列中最早的消息,队列的 MinOffset 和 MaxOffset 值会不断增加。Consumer offset update

消费者位点 (Consumer Offset)

Apache RocketMQ 采用发布-订阅模式。多个消费者组可以订阅同一个队列。在这种情况下,如果一个消费者在消费后删除了消息,其他消费者将无法再消费它。

为了防止这种情况发生,Apache RocketMQ 使用消费者位点来管理不同消费者的消息消费进度。Apache RocketMQ 不会在消息被消费后立即删除它,而是维护一个记录,记录消费者组已消费的最新消息,这被称为消费者位点(Consumer Offset)。

如果客户端重启,消费者能够根据服务器上保存的消费者位点继续处理消息。如果消费者位点过期并被删除,则服务器上保存的队列的 MinOffset 值将用作消费者位点。

提示

消费者位点保存在 Apache RocketMQ 服务器上,并从服务器恢复,与任何特定的消费者实例无关。因此,Apache RocketMQ 可以在不同的消费者之间恢复消费进度。

下图显示了消息队列中最小位点、最大位点和消费者位点之间的关系。Consumer progress

  • 消费者位点始终小于或等于最大位点。

    • 如果消息的生产和消费速度相同,且队列中没有未消费的消息,则消费者位点与最大位点相同。

    • 如果消费速度慢于生产速度,队列中就会存在未消费的消息。因此,消费者位点小于最大位点,其差值就是未消费消息的数量。

  • 通常,消费者位点大于或等于最小位点。如果消费者位点小于最小位点,消费者将无法消费消息。在这种情况下,服务器会向消费者恢复正确的消费者位点。

初始消费者位点

初始消费者位点是消费者组首次开始消费某个消息队列时,保存在服务器上的消费者位点。

当消费者首次从队列获取消息时,Apache RocketMQ 使用消息队列的最大位点作为初始消费者位点。换句话说,消费者从队列中的最新消息开始消费。

重置消费者位点

如果初始或当前的消费者位点与您的业务状态不符,您可以重置消费者位点以调整消费进度。

适用场景

  • 初始消费者位点不当:初始消费者位点默认为队列的最大位点,即客户端从最新消息开始消费。如果您需要消费更早的消息,可以将消费者位点重置为更早的位点。

  • 消费积压:如果消费者无法跟上消息生成的速度,可能会积压大量消息。如果积压的消息并非业务关键,您可以将消费者位点调整为更大的值,以跳过这些消息并减轻下游压力。

  • 业务回溯与修正处理:如果您希望重新消费因业务错误而导致处理不当的消息,可以将消费者位点设置为较小的值。

消费者位点重置功能

Apache RocketMQ 的消费者位点重置功能允许您:

  • 将消费者位点重置为消息队列中的任何位点。

  • 将消费者位点重置为特定的时间点。服务器会将消费者位点调整为最接近该时间点的位点。

限制

  • 重置消费者位点后,消费者将从新的位点开始消费消息。在回溯场景中,消费者会从历史消息(通常是冷数据)开始消费。这被称为“冷读取”,可能会给您的系统带来过重的负担。在进行此操作前,请评估风险与收益。建议对该权限实施严格的控制策略,以防止滥用和频繁重置。

  • Apache RocketMQ 仅允许重置可见消息的消费者位点。您无法重置处于调度状态或重试挂起状态的消息位点。有关更多信息,请参阅 延迟消息消费重试

版本兼容性

在不同版本的 Apache RocketMQ 中,服务器对初始消费者位点的定义有所不同。

  • 在 4.x 和 3.x 版本中,初始消费者位点定义为队列的消息状态。

  • 在 5.x 版本中,初始消费者位点是消费者开始接收消息时队列的最大位点。

因此,如果您是从较早版本升级而来,在启动客户端时必须注意初始消费者位点的变化。

使用说明

严格控制重置权限

重置消费者位点会给系统带来额外负担,并可能影响消息的读写。因此,建议您在执行此操作前评估风险与收益。