【问题标题】:Setting an initial "current-offset" and "lag" for new consumer groups for a given topic为给定主题的新消费者组设置初始“当前偏移量”和“滞后”
【发布时间】:2018-08-20 04:25:59
【问题描述】:

我正在开发一种产品,该产品可能会根据用户使用产品的方式添加/删除消费者组。

enable.auto.commit 在我们的产品中被关闭,而是在每次收到数据后提交偏移量。

我们最近实施了一项可以暂停/恢复产品的服务。 kafka 库(NodeJS)还没有可用的暂停/恢复功能,所以我最终根据消费者消费者组取消订阅/订阅主题,这似乎按我们的预期工作。

唯一的问题发生在添加新的消费者组时。首先,让我解释一下我看到的行为:

这里是消费者“group1”信息..

$ bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group philz-topic-group1

TOPIC                          PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG        CONSUMER-ID                                       HOST                           CLIENT-ID
philz-topic                    1          33              33              0          rdkafka-3ac4d56e-e94b-4365-9af7-04e485502b5d      /10.233.113.109                rdkafka
philz-topic                    4          34              34              0          rdkafka-d642805c-f5ea-4450-9cb0-3272fcbbffc9      /10.233.88.251                 rdkafka
philz-topic                    0          23              23              0          rdkafka-12cfca8b-fd61-4a68-bc5f-1946c8ef4eb1      /10.233.120.55                 rdkafka
philz-topic                    2          26              26              0          rdkafka-7561ca2a-9894-4a3d-83fe-d379bbe64fdf      /10.233.126.40                 rdkafka
philz-topic                    3          20              20              0          rdkafka-cd9d5ed6-7daa-4b75-8f39-6704c8d887ed      /10.233.119.133                rdkafka

这里是消费者“group2”的信息。消费者“group2”刚刚添加并完成了一项操作。因此,单个操作的 CURRENT-OFFSET 和 LAG 已更新。

$ bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group philz-topic-group2

TOPIC                          PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG        CONSUMER-ID                                       HOST                           CLIENT-ID
philz-topic                    3          -               20              -          rdkafka-b56306e1-b4b7-43fe-a604-ab7c12f70e9f      /10.233.119.133                rdkafka
philz-topic                    1          -               33              -          rdkafka-76c9a4d2-268b-4ebb-94a8-f1230c9bbfea      /10.233.113.109                rdkafka
philz-topic                    4          34              34              0          rdkafka-d412e574-8241-48c6-af26-c50be44eb51d      /10.233.126.40                 rdkafka
philz-topic                    0          -               23              -          rdkafka-33179a7d-cb9f-453a-83c6-e7e4780372b6      /10.233.88.251                 rdkafka
philz-topic                    2          -               26              -          rdkafka-77506e87-b666-4c92-82df-82071e2ff801      /10.233.120.55                 rdkafka

如果添加了新的消费者组并且没有完成任何操作,则上述命令不会显示有关消费者组的信息。

我目前面临的问题是,当发生暂停/恢复操作并且消费者组的所有分区都没有更新的 CURRENT-OFFSET 和 LAG 时,当取消订阅/暂停并完成操作时,分区应该有现在 LAG 为 1。但是,如果一个新的消费者组之前没有任何给定分区的 CURRENT-OFFSET 和 LAG,那么现在该信息将被跳过并且消费者组永远不会看到。

我的问题是,当创建一个新的消费者组时,我们可以更新组的 CURRENT-OFFSET 以匹配所有可用分区的 LOG-END-OFFSET 吗?

我对 Kafka 不是很熟悉,因此对此处行为的任何解释表示赞赏。

我的猜测是因为我们自己提交了offset(因为enable.auto.commit被关闭了),当一个操作发生时,我们能够看到新的消费者组的一些信息,但只能看到一个分区(刚刚那个接收到数据)显示并使用当前偏移量进行更新。

谢谢!

编辑:

另外,在我的示例中,每个消费者组有 5 个消费者和 5 个分区,因此每个分区应该有一个消费者

【问题讨论】:

  • 你也在设置auto.offset.reset吗?如果您开始一个新的消费者组,理想情况下您希望从主题的开头开始阅读,而不是最近的消息。
  • 我没有设置auto.offset.reset。那么您的意思是将其设置为“最早”吗?我要设置为最早的一个问题是用户可能会删除消费者组并重新添加消费者组的情况。 “最早”是否让每个消费者从第一个偏移量到最新偏移量读取一个分区?
  • 没有一个非常明显的 api 来删除消费者组,所以我不认为这应该是一个问题,但是是的,该设置应该确定从新组的最早可用偏移量开始
  • @cricket_007 太棒了!谢谢,刚刚试了一下,似乎将新组分区的偏移量设置为当前的日志结束偏移量。尝试一下,似乎解决了我的问题。如果您提交答案,我会接受 :) 我需要了解为我们产品中的其他用例设置此配置的含义,但现在,lgtm
  • 太棒了!我不熟悉 Node API,所以请随意发布您自己的答案

标签: node.js apache-kafka kafka-consumer-api


【解决方案1】:

感谢 cricket_007 提供了必要的 kafka 消费者选项

消费者选项auto.offset.reset 允许在实例化时自动设置消费者偏移量。通过将此选项的值设置为'earliest',它会将每个分区的当前偏移量设置为LOG-END-OFFSET

要使用节点库设置此选项,只需:

const consumer = new Kafka.KafkaConsumer(config, {
    'auto.offset.reset': 'earliest'
});

其中config 是您为消费者提供的键/值对配置,第二个参数是您的键/值对配置,用于创建默认主题配置。

配置是在消费者上设置的主题级配置,如下所述:https://github.com/edenhill/librdkafka/blob/0.11.1.x/CONFIGURATION.md

【讨论】:

  • 我认为这不是对原始问题的正确答案,因此将其删除为已接受的答案。原因是“最早”将偏移量设置为可能的最早偏移量,这可能导致重新使用给定主题的所有条目。但我真的在寻找一种方法来立即将偏移量设置为消费者对给定主题的创建和/或订阅的最新情况
  • 当前应用的修复需要在首次创建消费者并使用'auto.offset.reset': 'latest' 时在每个可用分区上发送和使用初始“​​空”消息。这样我们可以在取消订阅消费者之前设置偏移量。我将此视为临时修复/解决方法,因此不要发布为可接受的答案
猜你喜欢
  • 2016-10-31
  • 2019-05-01
  • 2019-04-04
  • 1970-01-01
  • 1970-01-01
  • 2021-11-10
  • 1970-01-01
  • 2018-10-27
  • 2018-07-29
相关资源
最近更新 更多