【问题标题】:Acknowledge within @KafkaListener-method without "losing" messages在@KafkaListener-method 中确认而不会“丢失”消息
【发布时间】:2017-05-18 14:14:35
【问题描述】:

我们目前基本上使用以下简化机制来确认消息:

  @KafkaListener(topics = "someTopic")
  public void listen(final String message, final Acknowledgment ack) {
    try {
        processMessage(message);
        ack.acknowledge();
    } catch (final IOException e) {
        // do not acknowledge here since we can temporarily not process the message
    } 

基本上,只要我们暂时无法处理消息(在 IOExceptions 的情况下),我们希望稍后再接收它。

但这不起作用,因为确认假定同一分区中的所有先前消息都已成功处理。在我们的 IOException 案例中,失败的消息会被跳过,但可能会被同一分区上具有更高索引的不同消息确认。

我们有一些解决这个问题的想法,但这意味着需要一些讨厌的解决方法来避免在 KafkaListener 方法中直接调用确认。我们的用例是一个非常具体的用例,还是更像是 spring kafka 用户会假设的“默认”行为?

这种问题有spring-kafka解决方案吗?或者你有一个“正确”解决这个问题的想法吗?

【问题讨论】:

    标签: spring apache-kafka spring-kafka


    【解决方案1】:

    这就是卡夫卡的工作方式;在继续下一条消息之前,您可以启用重试以尝试通过您的 IOException。

    您可以配置错误句柄以将失败消息发布到另一个主题以便稍后重播。或者,错误处理程序可以停止容器以阻止任何新的交付。

    当容器重新启动时,消息将被重播。

    【讨论】:

    • 感谢您的快速回答。我想我们用例的最佳解决方案是在移动到下一条消息之前使用重试。
    猜你喜欢
    • 1970-01-01
    • 2016-12-21
    • 2015-12-05
    • 2014-08-23
    • 1970-01-01
    • 2011-10-11
    • 2012-11-11
    • 2020-01-19
    • 1970-01-01
    相关资源
    最近更新 更多