|
|
1
1
让我们考虑一下您的消费者群名称是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
如果消费者收到消息n,提交它,然后在完全处理它之前崩溃,那么在默认情况下,消费者将认为该消息已处理。 请注意,消息仍然在代理上,因此可以重新使用它进行处理。但这要求应用程序中的某些逻辑不仅从上次提交的位置重新启动,而且还要检查以前的记录是否已成功处理。 如果应用程序处理消息通常需要很长时间,那么您可能希望切换到手动提交而不是自动提交。这样,您就能够更好地控制何时提交并避免这个问题。 |