【问题标题】:Spring Kafka acknowledgment settingsSpring Kafka 确认设置
【发布时间】:2022-01-13 06:28:30
【问题描述】:

我有一个 spring boot 应用程序,我们在消费者部分使用 spring Kafka 我已将 enable.auto.commit 设置为 false 并设置我的 listener ack-mode to manual_immediate

我有并发的消费者,所以在消费记录后我调用 acknowledgment.acknowledge() 但在这里我仍然面临重复问题问题,每当重新平衡发生时,其他消费者开始消费已经被一个消费者消费的相同消息。知道幕后发生了什么神奇的事情。

任何人都知道在使用 manual_immediate 时它是通过 commitSync 还是 commitAsync 提交消息?有没有办法我们可以改变行为以避免重复记录消息阅读。有没有办法我们可以在 Spring Kafka 中使用混合模型

在 Spring Boot Kafka 中,有一种方法可以让我们在重新平衡发生时看到我可以记录它。

如果我们想为某些测试目的创建再平衡?

【问题讨论】:

    标签: java spring-boot apache-kafka spring-kafka


    【解决方案1】:

    只要你在监听线程上调用acknowledge,它就会默认使用commitSync();使用 syncCommits 容器属性来使用异步提交。

    如果您在不同的线程上调用它,提交将排队等待由消费者线程尽快处理。

    如果由于您的侦听器处理poll() 收到的记录花费的时间太长而发生强制重新平衡,则无法避免重复。

    您可以增加max.poll.interval.ms和/或减少max.poll.records,以确保及时处理记录。

    您可以将ConsumerRebalanceListener 添加到容器属性以记录重新平衡。

    max.poll.interval.ms 减小到一个较小的值以在测试中重现。

    【讨论】:

    • 嘿,max.poll.interval.ms 的作用是什么?这是否意味着默认为 5 分钟,所以如果消耗消息的成本超过 5 分钟,它会导致重新平衡,对吗?
    • 是的,没错。
    • 5 min to consume here是指处理批处理中的所有消息还是单个消息?
    • 上次投票的所有消息 - 正如我所说,您可以使用 max.poll.records 减少批量大小。
    【解决方案2】:

    首先,无论 ack 模式如何,都不能保证一条消息只被消费一次。例如,在消费消息和提交时间偏移之间可能会发生重新平衡,导致 Kafka 再次将消息传递给新分配的消费者。对重复消息幂等是应用程序的责任。

    为了监听重新平衡事件,需要实现ConsumerRebalanceListener。您可以将此实现插入 Spring 的自动配置的 ConcurrentKafkaListenerContainerFactory 实例。 here 已经回答了有关如何完成此操作的更详细说明。

    如果您希望创建强制重新平衡以进行测试,您可以通过杀死希望超过 1 个现有消费者中的一个来实现。如果使用 spring-kafka,您可以通过使用 @AutowiredKafkaListenerEndpointRegistry 实例并杀死/休息/(重新)启动任何消费者来做到这一点。这方面应该做的事情:

    @Autowired
    KafkaListenerEndpointRegistry registry;
    
    public void myTest() {
        Collection<MessageListenerContainer> containers = registry.getAllListenerContainers()
        containers.get(0).stop()
    }
    

    【讨论】:

    • thankyou @Akhil Bojedla 如何控制 Kafka 将消息再次传递给新分配的消费者的问题,您能否提供一个小的编码示例 ConsumerRebalanceListener 以便我可以继续使用它,这对我最有帮助
    • 如前所述,这应该在您的消费者内部处理。这是事件驱动系统的常见模式。不幸的是,很难给出代码示例,因为制作幂等系统是针对个别用例的。但是,您可以查找一些使应用程序具有幂等性的常见高级模式。我敢肯定,对幂等消费者模式的快速谷歌搜索应该会找到一些写得很好的文章。
    • 嘿,max.poll.interval.ms 的作用是什么?这是否意味着默认为 5 分钟,所以如果消耗消息的成本超过 5 分钟,它会导致重新平衡,对吗?
    猜你喜欢
    • 2020-06-29
    • 1970-01-01
    • 2021-09-28
    • 2017-04-29
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-05-08
    • 2020-05-19
    相关资源
    最近更新 更多