【问题标题】:What determines Kafka consumer offset?是什么决定了 Kafka 消费者偏移量?
【发布时间】:2015-11-30 03:39:19
【问题描述】:

我对卡夫卡比较陌生。我已经对它进行了一些试验,但是关于消费者补偿,我还不清楚一些事情。据我目前了解,当消费者启动时,它将开始读取的偏移量由配置设置auto.offset.reset 确定(如果我错了,请纠正我)。

现在假设主题中有 10 条消息(偏移量 0 到 9),而消费者恰好在它关闭之前(或在我杀死消费者之前)消费了其中的 5 条。然后说我重新启动该消费者进程。我的问题是:

  1. 如果auto.offset.reset 设置为earliest,它是否总是从偏移量 0 开始消费?

  2. 如果auto.offset.reset 设置为latest,它会从偏移量5 开始消费吗?

  3. 这种情况下的行为是否总是确定性的?

如果我的问题中有任何不清楚的地方,请随时发表评论。

【问题讨论】:

    标签: java apache-kafka kafka-consumer-api distributed-computing


    【解决方案1】:

    它比你描述的要复杂一点。
    auto.offset.reset config 仅在您的消费者组没有在某处提交有效偏移量时才会生效(现在支持的 2 个偏移量存储是 Kafka 和 Zookeeper),而且它还取决于您使用哪种消费者。

    如果您使用高级 java 消费者,那么想象以下场景:

    1. 您在消费者组 group1 中有一个消费者,该消费者已消费 5 条消息并死亡。下次你启动这个消费者时,它甚至不会使用 auto.offset.reset 配置,而是会从它死去的地方继续,因为它只会从偏移存储(我提到的 Kafka 或 ZK)中获取存储的偏移。

    2. 您在一个主题中有消息(如您​​所描述的),并且您在一个新的消费者组 group2 中启动了一个消费者。任何地方都没有存储偏移量,这次auto.offset.reset 配置将决定是从主题的开头(earliest)还是从主题的结尾(latest)开始

    影响与earliestlatest 配置对应的偏移值的另一件事是日志保留策略。假设您有一个保留时间配置为 1 小时的主题。您生成 5 条消息,然后一个小时后您又发布了 5 条消息。 latest 偏移量仍将与上一个示例中的相同,但 earliest 不能是 0,因为 Kafka 已经删除了这些消息,因此最早可用的偏移量将是 5

    上面提到的一切都与SimpleConsumer无关,每次运行它都会决定从哪里开始使用auto.offset.reset配置。

    如果您使用 Kafka 版本早于 0.9,则必须将 earliestlatest 替换为 smallestlargest

    【讨论】:

    • 非常感谢您的回答。那么对于高级消费者来说,一旦消费者提交了某些事情(在 ZK 或 Kafka 中),auto.offset.reset 此后就没有任何意义了吗?该设置的唯一意义是什么时候没有提交(理想情况下是在消费者第一次启动时)?
    • 和你描述的完全一样
    • @serejja 你好 - 如果我总是每组有 1 个消费者,并且你的答案的场景#1 发生在我身上,怎么样?会一样吗?
    • @ha9u63ar 不太明白你的问题。如果您在同一个组中重新启动您的消费者,那么是的,它不会使用 auto.offset.reset 并从提交的偏移量继续。如果你总是使用不同的消费者组(比如启动消费者时生成它),那么消费者将永远尊重auto.offset.reset
    • @serejja 是的,这对我不起作用。你能看看this - 这是我的问题
    【解决方案2】:

    只是一个更新:从 Kafka 0.9 及以后,Kafka 使用新的 Java 版本的消费者,并且 auto.offset.reset 参数名称已更改;来自手册:

    当 Kafka 中没有初始偏移量或当前 服务器上不再存在偏移量(例如,因为该数据 已被删除):

    最早的:自动将偏移量重置为最早的偏移量

    最新的:自动将偏移量重置为最新的偏移量

    none:如果没有找到先前的偏移量,则向消费者抛出异常 针对消费者群体

    其他:向消费者抛出异常。

    我在检查接受的答案后花了一些时间找到这个,所以我认为发布它可能对社区有用。

    【讨论】:

    • 接受的答案是用新名字写的——这个答案没有提供任何独特的东西,不是吗? (如果在撰写本文时没有 90 票,我建议删除它;))
    • 令人惊讶的是很多人发现它很有用。
    • 我同意一个答案不会完全偶然地获得那么多赞成票。但是关于原始答案的观点不再是 AFAICT,所以我想不出我现在要投票的理由吗? (在登陆这里之前,我也看过手册的特定部分)。另外:this answer 在这个领域也很有用
    【解决方案3】:

    还有 offsets.retention.minutes。如果自上次提交以来的时间是 > offsets.retention.minutes,那么 auto.offset.reset 也会生效

    【讨论】:

    • 这似乎与日志保留无关?偏移保留应该基于日志保留吗?
    • @mike01010 没错。它应该基于日志保留,这是工单中建议的解决方案之一。 Prolong default value of offsets.retention.minutes to be at least twice larger than log.retention.hours.issues.apache.org/jira/browse/KAFKA-3806
    • 这个答案让我害怕了一段时间,直到我检查了the documentationoffsets.retention.minutes在消费者组失去所有消费者(即变为空)后,它的偏移量将为此保留在被丢弃之前的保留期。 对于独立消费者(使用手动分配),偏移量将在最后一次提交时间加上此保留期之后过期。 (这是给Kafka 2.3
    猜你喜欢
    • 1970-01-01
    • 2017-07-22
    • 2019-05-01
    • 2018-02-03
    • 1970-01-01
    • 2016-03-05
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多