跳至主内容

使用 RabbitMQ Streams 进行偏移量跟踪

·9 分钟阅读

RabbitMQ Streams 为消费者提供服务器端偏移量跟踪。此功能允许消耗应用程序在上次中断的地方恢复消耗。本文介绍偏移量跟踪的语义以及它在 Streams Java 客户端中的实现方式。

什么是偏移量跟踪?

偏移量跟踪是消费应用程序存储其在流(stream)中位置的过程。如果一个消费应用程序在流中的位置是 1,000,这意味着它已经处理了该位置之前的所有消息。如果应用程序停止并重新启动,它将从其最后存储的位置(在本例中为 1,001)之后重新连接到流,并从那里重新开始消费。

我们使用偏移量 (offset) 的概念来指定流中的确切位置。

首先,什么是偏移量?

如果我们把流建模为一个数组,其中每个元素都是一条消息,那么偏移量就是给定消息在该数组中的索引。

A stream can be represented as an array. The offset is the index of an element in the array.
流可以表示为一个数组。偏移量是数组中元素的索引。

这种思维模型足以理解我们这里的偏移量跟踪,但它并不完全准确。我们可以说,这个大数组中的元素并不总是消息,它可能是流内务管理所需的一些信息。一个结果是,2 条连续的消息可能并不总是有 2 个连续的偏移量值,但这在本文的背景下并不重要。

流中不同的偏移量规范

绝对偏移量没有任何意义,它只是技术性的。因此,当应用程序第一次想要连接到流时,它不太可能使用偏移量,它更倾向于使用更高级的概念,例如流的开始 (beginning)结束 (end),甚至是流中的某个时间点 (point in time)

幸运的是,RabbitMQ Streams 除了绝对偏移量外,还支持不同的偏移量规范firstlastnexttimestamp

流的“末尾”有两种偏移量规范:next 表示要写入的下一个偏移量。如果消费者在 next 处连接到流且没有人发布消息,消费者将不会收到任何内容。当有新消息进入时,它将开始接收。

last 意味着“从最后一组消息开始”。chunk 是 RabbitMQ Streams 中“一批消息”的术语,因为出于性能原因,消息是分批处理的。

下图显示了流中的偏移量规范。

Different offset specifications in a stream with 2 chunks.
包含 2 个 chunk 的流中的不同偏移量规范。

偏移量跟踪的不同步骤

因此,应用程序通常会在第一次启动时指定像 firstnext 这样的偏移量,处理消息,并定期(例如每 10,000 条消息)在服务器上存储一个偏移量值。

应用程序可能会因为某种原因(例如升级)而停止。当它重新启动时,它将检索其最后存储的偏移量,并在该位置之后重新开始消费。

RabbitMQ Stream 协议提供了用于在给定流上为给定应用程序存储和查找偏移量的命令,因此偏移量跟踪主要是客户端关注的问题。

偏移量跟踪的应用程序要求

对于想要在流中使用偏移量跟踪的应用程序,有一个简单的要求:它必须使用在应用程序重启后保持不变的跟踪引用 (tracking reference)

指定此引用的方式取决于客户端库,更一般地说,客户端库可能具有处理偏移量跟踪的不同特性或方法。点击此处阅读更多关于流 Java 客户端中偏移量跟踪的信息。

服务器端偏移量跟踪的阴暗面

RabbitMQ Streams 中的服务器端偏移量跟踪非常简洁,但应谨慎使用。偏移量跟踪信息存储在流日志中。还记得流的数组表示吗?想象一下,数组元素大部分是消息,但其中一些包含偏移量跟踪信息,例如“应用程序 'my-application' 的偏移量 1000”。

因此,每当客户端要求为给定消费者存储偏移量时,都会在日志中创建一个小条目。你肯定不希望创建过多的此类条目,这就是为什么我们建议每几千条消息存储一次偏移量,并避免为每一条消息都存储偏移量的消费者。同样,对于服务器端偏移量跟踪,请务必谨慎和合理。

还要记住,服务器端偏移量跟踪只是一种便利设施,绝不是允许消费者从中断处重启的唯一解决方案。想象一下,消息处理是事务中的数据库操作。将偏移量作为同一事务的一部分存储在数据库中是一个好主意,因为它确保了数据更改和偏移量存储以原子方式发生。

服务器端偏移量跟踪实战

现在让我们看看如何使用流 Java 客户端设置偏移量跟踪。

跟踪消费者

这里有一些用于启动带有偏移量跟踪的消费者的代码

AtomicInteger messageConsumed = new AtomicInteger(0);
Consumer consumer = environment.consumerBuilder()
.stream("offset-tracking-stream") // the stream to consume from
.offset(OffsetSpecification.first()) // start consuming at the beginning
.name("my-application") // the name (reference) of the consumer
.manualTrackingStrategy() // tracking is done in application code
.builder()
.messageHandler((context, message) -> {
// ... message processing ...

// condition to store the offset: every 10,000 messages
if (messageConsumed.incrementAndGet() % 10_000 == 0) {
context.storeOffset(); // store the message offset
}
// ...
})
.build();

如果你需要回顾流 Java 客户端 API,可以阅读 RabbitMQ Streams 首个应用程序

以下是此代码片段中的关键点

  • 消费者必须有一个名称才能启用偏移量跟踪。它在存储偏移量值时使用其名称作为跟踪引用。
  • 消费者使用 OffsetSpecification.first() 从流的开头开始消费。当消费者重新启动且存在为其存储的偏移量值时,此规范将被忽略。
  • 应用程序代码使用 Context#storeOffset() 显式处理跟踪。Consumer#store(long) 方法是存储偏移量的另一种可能性。

你可以看到偏移量跟踪如何依赖于客户端库。如果你不想在代码中处理偏移量跟踪,流 Java 客户端也提供了一种自动偏移量跟踪策略

如你所见,当理解了需求和语义后,偏移量跟踪非常简单。现在让我们运行该示例。

设置示例项目

运行示例需要安装 Docker、Git 和 Java 8 或更高版本。您可以使用以下命令启动代理:

docker run -it --rm --name rabbitmq -p 5552:5552 \
-e RABBITMQ_SERVER_ADDITIONAL_ERL_ARGS='-rabbitmq_stream advertised_host localhost' \
rabbitmq:3.9

然后您需要启用 Streams 插件:

docker exec rabbitmq rabbitmq-plugins enable rabbitmq_stream

代码托管在 GitHub 上。以下是如何克隆存储库的方法

git clone https://github.com/acogoluegnes/rabbitmq-streams-blog-posts.git
cd rabbitmq-streams-blog-posts

在示例项目中,我们将

  • 发布第一波消息
  • 启动消费者,它将消费这些消息,定期存储偏移量,并在第一波消息结束时停止(多亏了毒药消息)
  • 发布第二波消息
  • 再次启动消费者,并确保它从中断处(第一波消息结束时)重新开始,并且只消费来自第二波的消息

发布第一波消息

使用以下命令发布第一波消息

./mvnw -q compile exec:java -Dexec.mainClass='com.rabbitmq.stream.OffsetTracking$PublishFirstWave'

你应该在控制台上得到以下内容

Connecting...
Connected
Creating stream...
Stream created
Creating producer...
Producer created
Sending 500,000 messages
Messages sent, waiting for confirmation...
All messages confirmed? yes (928 ms)
Closing environment...
Environment closed

该程序发布了 500,000 条消息,它们具有相同的 first wave 主体,除了最后一条消息是用于停止消费的毒药消息。

消费第一波消息

启动消费者

./mvnw -q compile exec:java -Dexec.mainClass='com.rabbitmq.stream.OffsetTracking$Consume'

你得到

Connecting...
Connected
Start consumer...
Consumed 500,000 messages in 364 ms (bodies: poison, first wave)
Closing environment...
Environment closed

消费者的表现符合预期:它从头开始,因为它消费了 500,000 条消息。消费者还将它获取的所有主体记录在一个集合中,它获取了 first wave 消息和 poison 消息,很好。

发布第二波消息

让我们发布另一波消息

./mvnw -q compile exec:java -Dexec.mainClass='com.rabbitmq.stream.OffsetTracking$PublishSecondWave'

这次我们得到

Connecting...
Connected
Creating producer...
Producer created
Sending 100,000 messages
Messages sent, waiting for confirmation...
All messages confirmed? yes (312 ms)
Closing environment...
Environment closed

好的,流中又多了 100,000 条消息。

消费第二波消息

让我们重启消费者

./mvnw -q compile exec:java -Dexec.mainClass='com.rabbitmq.stream.OffsetTracking$Consume'

我们得到

Connecting...
Connected
Start consumer...
Consumed 100,000 messages in 215 ms (bodies: poison, second wave)
Closing environment...
Environment closed

太棒了!即使消费者本应从流的开头开始,它也检测到存在存储的偏移量,因此它使用它来正确重启。它只消费了第二波消息。

总结

RabbitMQ Streams 为消费者提供服务器端偏移量跟踪。消费者实例必须具有一个在重启后保持不变的名称(或跟踪引用)。客户端应用程序定期存储偏移量,具体方式取决于客户端库。

对于需要从中断处重新启动的应用程序,偏移量跟踪是必不可少的。

流 Java 客户端对偏移量跟踪提供了极好的支持,其文档对此进行了深入介绍。

© . This site is unofficial and not affiliated with VMware.