【问题标题】:Backoff strategy for apache kafka consumerapache kafka消费者的退避策略
【发布时间】:2020-01-02 14:46:31
【问题描述】:

骆驼使用kafka组件时,从kafka消费时有两种重试方式:

  • 内存重试,使用骆驼路线上的通用错误处理。但问题是,在重试时,消费者停止轮询代理,如果达到 max.poll.interval.ms,Kafka 代理认为消费者不健康,并将其从消费者组中删除:

org.apache.kafka.clients.consumer.internals.AbstractCoordinator | [消费者clientId=consumer-1, groupId=2862121d-ddc9-4111-a96a-41ba376c0143] 该成员将离开 组,因为消费者轮询超时已过期。这意味着 后续调用 poll() 之间的时间比配置的长 max.poll.interval.ms,这通常意味着轮询循环是 花太多时间处理消息。你可以解决这个问题 通过增加 max.poll.interval.ms 或减少最大值 使用 max.poll.records 在 poll() 中返回的批次大小。

  • 使用参数 breakOnFirstError 在每次重试时轮询。偏移量没有更新,我们不断从代理轮询相同的消息。问题是我找不到定义退避策略的方法,并且重试的次数太频繁了。

您知道如何为第二种方法定义退避策略吗?

【问题讨论】:

    标签: apache-kafka apache-camel


    【解决方案1】:

    我不熟悉 Apache Camel,但是如果你能够修改消费者参数和轮询循环,那么第二种方法在这里是正确的,它是 Kafka 重试方式 - 不要提交偏移量,所以接下来轮询循环迭代将再次使用该消息。

    进一步的策略取决于您在处理故障时究竟需要什么:

    • 您希望重试最终成功吗?然后,为了避免向同一条消息发送垃圾邮件,您可以调整消费者使用 max.poll.interval.ms 配置参数从 Kafka 轮询消息的时间间隔。更多详情here

    • 您要重试一定次数然后继续下一条消息吗?在这种情况下,您需要在轮询循环中手动实现重试计数器。一旦达到一定的重试次数 - 您只需进一步推动消费者:

      final TopicPartition topicPartition = new TopicPartition(topic, partition); consumer.seek(topicPartition, consumer.position(topicPartition) + 1);

    【讨论】:

    • 感谢您的回答。我希望重试最终会成功。实际上,它发生在服务器暂时关闭时。我想定义一个指数退避策略,因为大多数时候服务器都在运行,我想尽可能快地进行轮询,但是每当服务器关闭时,我想避免垃圾邮件。
    猜你喜欢
    • 2015-06-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-08-18
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多