【问题标题】:How Do I know whether my record has been Manually committed using Spring Kafka我如何知道我的记录是否已使用 Spring Kafka 手动提交
【发布时间】:2017-04-29 22:53:51
【问题描述】:

我想知道当 AckMode 在 spring kafka 中设置为 MANUAL 时提交是如何工作的。

下面是我在KafkaConfig中设置的属性 containerProperties.setAckMode(AbstractMessageListenerContainer.AckMode.MANUAL);

listener 代码

@KafkaListener(id="POC", topics = "TestTopic", group = "TestGroup")
    public void listen(ConsumerRecord<String,KafkaPayload> record, Acknowledgment acknowledgment) {
        countDownLatch.countDown();     
        acknowledgment.acknowledge();
}

我正在按照 spring kafka 文档执行acknowledgement,但这仅意味着我的消息被标记为 sent 而不是 consumed(这是我的理解)。

  1. 在这种情况下,我应该调用 commitSync() 方法吗? 如果,我从哪里调用它,因为我需要获得对KafkaConsumer 的引用。 如果否,它是如何在内部工作的,我可以跟踪它吗?

  2. 是否有 commitId 或某个返回值? 我的想法是知道是否消费了特定的消费者记录。 我想存储该值以用于内部跟踪。

  3. kafka 是否在内部维护消费者记录中的任何状态,例如(AcknowledgedCommittedNot Committed),这有助于分类。

这真的可以帮助我区分有多少记录被消耗,有多少是待处理的以及它们的状态。

【问题讨论】:

    标签: apache-kafka kafka-consumer-api spring-kafka


    【解决方案1】:

    我可以回答第一个问题。其余的一切看起来都像是 Apache Kafka 的直接故事。

    由于我们不能从我们想要的地方执行commit,而只能从执行consumer.poll()的同一个线程中执行,我们将所有提交请求存储在内部KafkaMessageListenerContainer队列中,并在执行this.consumer.poll()之前的主要消费者循环。

    即使您使用MANUAL_IMMEDIATE,真正的consumer.commitSync() 也是在与您的acknowledgment.acknowledge() 不同的线程上执行的。

    OTOH 在Consumer 中的API 看起来像:

    public void commitSync(Map<TopicPartition, OffsetAndMetadata> offsets);
    

    所以,没有任何commitId 钩子可以解决。

    我认为 Apache Kafka 中没有像 Not Committed 或其他任何东西这样的概念。数据始终存在于主题日志中,并且在特定的管理操作或压缩配置之前不会从那里删除。

    我认为commit offset 功能与consumer group 目的完全相关,根据我们拥有的JavaDocs:

    * This commits offsets to Kafka. The offsets committed using this API will be used on the first fetch after every
    * rebalance and also on startup. As such, if you need to store offsets in anything other than Kafka, this API
    * should not be used. The committed offset should be the next message your application will consume,
    * i.e. lastProcessedMessageOffset + 1.
    

    因此,当您的消费者死亡时,它将从其组的上次提交的偏移量重新启动。不同的组可能会读取相同的数据,但来自其他偏移量。我认为这绝对是为什么他们的 API 没有提供任何挂钩到实际状态的原因。就是没有这样的人!

    【讨论】:

    • 感谢 Artem 的详细回复。所以,我从上面的答案中了解到,当我们调用acknowledgment.acknowledge() 时,consumer.commitSync() 发生在不同的线程上。因此,这确保了提交正在发生。 'commitAsync()` 的行为是否也相同?我希望 AckMode.MANUAL 也一样,因为我没有使用 MANUAL_IMMEDIATE
    • 是的,这是真的。任何Consumer 操作都发生在同一个线程上。 MANUAL_IMMEDIATEMANUAL 的不同之处仅在于立即 consumer.wakeup() 打破当前投票。
    猜你喜欢
    • 1970-01-01
    • 2010-09-16
    • 1970-01-01
    • 2013-11-30
    • 2020-01-08
    • 1970-01-01
    • 2017-09-10
    • 1970-01-01
    • 2017-02-06
    相关资源
    最近更新 更多