【问题标题】:Rereading message from Kafka topic by refusing acknowledgement通过拒绝确认重读来自 Kafka 主题的消息
【发布时间】:2016-01-20 20:50:48
【问题描述】:

我正在使用spring-integration-kafka 实现具有自定义确认机制的 Kafka 消费者。

使用了来自this example 的代码。

我想要实现的是,当抛出异常时,不应将确认发送回 Kafka(即不应执行偏移提交),因此下一个 fromKafka.receive(10000) 方法调用将返回与上一个。

但是我遇到了一个问题:即使没有向 Kafka 发送确认,消费者也以某种方式知道下一条消息的偏移量并继续读取新消息,尽管偏移量主题中的偏移量值保持不变.

如果出现一些失败,如何让消费者重读消息?

【问题讨论】:

    标签: spring-integration apache-kafka


    【解决方案1】:

    目前不支持重新获取失败的消息

    您可以做的一件事是在消息驱动的适配器下游添加重试(例如,使用请求处理程序重试建议)。

    通过不确认,消息将在重新启动后传递,但不会在当前实例化期间传递。

    由于消息已预取到适配器中,您可以做的一件事是检测故障、停止适配器、耗尽预取的消息并重新启动。

    您可以注入自定义 ErrorHandler 来停止适配器并向下游流发出信号,表明它应该忽略排空消息。

    编辑

    现在有一个SeekToCurrentErrorHandler

    【讨论】:

    • 谢谢,我会尝试使用ErrorHandler 的方式。但无论如何,适配器如何知道下一条消息的偏移量?它是否在内部存储当前偏移量?
    • @gary 在最新的 spring kafka 中,有没有办法在我们手动 ack 时重新获取失败的消息而不是 fwd?只是为了避免重新启动适配器。我尝试了 Seek,它不会回到失败的偏移量。使用 spring-kafka 2.1.4 与集成 kafka 3.0.3
    • 是的,现在有一个SeekToCurrentErrorHandler。当前版本是 2.1.8。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-07-09
    • 1970-01-01
    • 2018-07-18
    • 1970-01-01
    • 2019-11-15
    • 1970-01-01
    • 2018-01-09
    相关资源
    最近更新 更多