【问题标题】:Kafka consumer crashing卡夫卡消费者崩溃
【发布时间】:2020-06-11 07:26:16
【问题描述】:

我的应用程序使用来自Kafka 0.9 的数据。

目前,我们获取消息、处理和提交。在这个过程中,如果消费者因为处理时间较长而无法发送心跳,消费者协调器会认为消费者已经死亡,并重新平衡分区。因此,总体而言,这会导致消费者数量减少。在一段时间后,我们的数据处理会停止。

我应该如何处理这种应用程序失败?

如果协调器发现一个新的消费者已经死了,有没有办法让消费者保持活力或跨越一个新的消费者?

【问题讨论】:

  • 您的“死”消费者是否尝试重新连接?我记得,在提交偏移量时它应该会失败。
  • 最小化心跳,让协调者知道消费者还活着。

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


【解决方案1】:

你需要增加max.poll.interval.ms的值 这为消费者在获取更多记录之前可以空闲的时间量设置了上限。 其他消费者心跳配置请参考this

【讨论】:

    【解决方案2】:

    一个明显的解决方案是分析处理消息的代码并尝试在其中减少阻塞工作。例如,同一线程上的 HTTP 或数据库调用

    日志应该告诉您减少 max.poll.records(在轮询之间做更少的总处理)或增加 max.poll.interval.ms(增加轮询之间的等待时间)

    消费者心跳还有其他设置,但从这些开始

    【讨论】:

      【解决方案3】:

      0.9 非常非常老了;这方面有很多改进。见KIP-62

      新消费者的一个常见问题是其单线程设计和需要通过向协调器发送心跳来保持活跃度。我们建议用户从与消费者轮询循环相同的线程进行消息处理和分区初始化/清理,但如果这花费的时间超过配置的会话超时时间,则消费者将从组中删除,并将其分区分配给其他成员。 ...

      如果您使用 Spring for Apache Kafka,您应该至少使用 1.3.10 版本;不再支持早期版本。

      当前版本是 2.4.3。

      如果您使用的是当前的spring-kafka版本,则需要减少max.poll.records或增加max.poll.interval.ms

      【讨论】:

        猜你喜欢
        • 2017-01-07
        • 2021-03-11
        • 2019-07-03
        • 2018-05-05
        • 2021-08-22
        • 1970-01-01
        • 1970-01-01
        • 2020-10-28
        • 2015-12-18
        相关资源
        最近更新 更多