【发布时间】:2020-07-29 09:52:36
【问题描述】:
我们的一个 Kafka Streams 应用程序的 StreamThread 消费者在生成以下日志消息后进入僵尸状态:
[Consumer clientId=notification-processor-db9aa8a3-6c3b-453b-b8c8-106bf2fa257d-StreamThread-1-consumer, groupId=notification-processor] 成员notification-processor-db9aa8a3-6c3b-453b-b8c8-106bf2fa257d-StreamThread-由于消费者轮询超时,1-consumer-b2b9eac3-c374-43e2-bbc3-d9ee514a3c16 向协调器 ****:9092 (id: 2147483646 rack: null) 发送 LeaveGroup 请求已过期。这意味着后续调用 poll() 之间的时间比配置的 max.poll.interval.ms 长,这通常意味着轮询循环花费了太多时间来处理消息。您可以通过增加 max.poll.interval.ms 或通过使用 max.poll.records 减少 poll() 返回的批次的最大大小来解决此问题。
StreamThread 的 Kafka Consumer 似乎已经离开了消费组,但 Kafka Streams App 仍然处于 RUNNING 状态,没有消费任何新记录。
我想检测到 Kafka Streams 应用程序已进入这种僵尸状态,以便可以将其关闭并替换为新实例。通常,我们通过 Kubernetes 运行状况检查来验证 Kafka Streams 应用程序是否处于 RUNNING 或 REPARTITIONING 状态,但这不适用于这种情况。
因此我有两个问题:
- 当 Kafka Streams 应用程序没有活动的消费者时,它是否会保持在 RUNNING 状态?如果是:为什么?
- 我们如何检测(以编程方式/通过指标)Kafka Streams 应用已进入没有活跃消费者的僵尸状态?
【问题讨论】:
标签: java apache-kafka apache-kafka-streams confluent-platform