代码之家  ›  专栏  ›  技术社区  ›  mmdc

如果使用者保存消息的时间比自动提交间隔时间长,Kafka是否会丢失消息?

  •  1
  • mmdc  · 技术社区  · 8 年前

    假设自动提交间隔时间为30秒,则由于某些原因,使用者无法处理消息并将其保持超过30秒,然后崩溃。自动提交偏移机制是否在使用者崩溃之前就提交了这个偏移?

    如果我的假设是正确的,那么消息将丢失,因为它的偏移量已提交,但消息本身尚未被处理?

    2 回复  |  直到 8 年前
        1
  •  1
  •   Indraneel Bende    8 年前

    让我们考虑一下您的消费者群名称是test,并且您在消费者群中只有一个消费者。

    启用自动提交时,仅在poll()调用期间和关闭使用者期间提交偏移量。

    例如-auto.commit.interval.ms是5秒,每次调用poll()需要7秒。每次调用poll()时,它都会检查自动提交间隔是否已经过,如果已经过了,就像上面的例子一样,它将提交偏移量。

    补偿也在消费者结算时承诺。

    从文档中-

    “关闭使用者,等待30秒的默认超时以进行任何必要的清理。如果启用了自动提交,如果可能,这将在默认超时内提交当前偏移量”。

    你可以在这里了解更多-

    https://kafka.apache.org/10/javadoc/index.html?org/apache/kafka/clients/consumer/KafkaConsumer.html

    现在,关于您的问题,如果poll()不再被调用,或者使用者没有关闭,那么它将不会提交偏移量。

        2
  •  0
  •   Mickael Maison    8 年前

    如果消费者收到消息n,提交它,然后在完全处理它之前崩溃,那么在默认情况下,消费者将认为该消息已处理。

    请注意,消息仍然在代理上,因此可以重新使用它进行处理。但这要求应用程序中的某些逻辑不仅从上次提交的位置重新启动,而且还要检查以前的记录是否已成功处理。

    如果应用程序处理消息通常需要很长时间,那么您可能希望切换到手动提交而不是自动提交。这样,您就能够更好地控制何时提交并避免这个问题。

    推荐文章