【问题标题】:How does Kafka provides next batch of records to poll when commitAsync gets failed in committing offset当 commitAsync 提交偏移失败时,Kafka 如何提供下一批记录进行轮询
【发布时间】:2020-08-05 18:55:18
【问题描述】:

我有一个关于 Kafka 消费者使用记录的用例。 例如, 我有 1 个主题有 1 个分区。目前,它有 10 条记录,在消耗前 10 条记录时,另外 10 条记录被写入分区。

  1. myConsumer 第一次轮询并返回前 10 条记录,例如 0 - 9 条记录。
  2. 成功处理了所有记录。
  3. 它向 Kafka 调用了 commitAsync() 以提交最后一个偏移量。
  4. 提交响应正在处理中。这可能是成功的,也可能是失败的。
  5. 但是,由于是异步模式,它会继续轮询下一批。
  6. 现在,Kafka 或消费者民意调查如何知道它必须从第 10 个位置读取?因为 commitAsync 请求尚未完成。

请帮助我理解这个概念。

【问题讨论】:

    标签: apache-kafka kafka-consumer-api


    【解决方案1】:

    Commit Offset 告诉 broker 消费者已经成功处理了相应的消息。消费者自己会知道它的进度(除了消费者的开始,它从代理那里获得最后提交的偏移量)。

    在您的描述中的第 5 步,提交偏移量正在进行中。所以:

    • 经纪人知道已处理 0-9 条记录
    • 消费者本身已经阅读了消息,因此它知道自己已经阅读了 0-9 条消息。所以它会知道接下来从 10th 开始阅读。

    可能的场景

    1. 假设提交失败 (0-9)。你的下一批,比如 (10-15) 被成功处理和提交,那么就没有造成任何伤害。因为我们向代理标记到 15 的处理已完成。
    2. 假设提交失败 (0-9)。您的下一批 (10-15) 已处理,并且在提交之前,消费者已关闭。当您的消费者重新启动时,它会从代理获取其状态(这两个批次都没有提交)。因此它将从第 0 条消息开始读取。

    您也可以提出其他几种方案。我想底线是,当您的消费者因任何原因重新启动并且它已从 kafka 代理获得其最后处理的偏移量时,提交的重要性就会显现出来。

    【讨论】:

    • 谢谢你,@Rishabh Sharma。我明白了。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-12-10
    • 2016-05-12
    • 2022-10-04
    • 1970-01-01
    • 2018-06-08
    • 2021-03-17
    相关资源
    最近更新 更多