【问题标题】:Kafka Consumer error: Marking coordinator dead卡夫卡消费者错误:标记协调员死亡
【发布时间】:2018-05-27 00:06:01
【问题描述】:

我在 Kafka 0.10.0.1 集群中有一个包含 10 个分区的主题。我有一个产生多个消费者线程的应用程序。对于这个主题,我将产生 5 个线程。在我的应用程序日志中多次看到此条目

INFO :: AbstractCoordinator:600 - Marking the coordinator x.x.x.x:9092
(id:2147483646 rack: null) dead for group notifications-consumer

然后有几个条目说(Re-)joining group notifications-consumer. 之后我还看到一个警告说

Auto commit failed for group notifications-consumer: Commit cannot be completed since
the group has already rebalanced and assigned the partitions to another member. This means
that the time between subsequent calls to poll() was longer than the configured
max.poll.interval.ms, which typically implies that the poll loop is spending too much time 
message processing. You can address this either by increasing the session timeout
or by reducing the maximum size of batches returned by poll() with max.poll.records.

现在我已经像这样调整了我的消费者配置

props.put("max.poll.records", 200);
props.put("heartbeat.interval.ms", 20000);
props.put("session.timeout.ms", 60000);

因此,即使在正确调整配置后,我仍然会收到此错误。在重新平衡期间,我们的应用程序完全没有响应。请帮忙。

【问题讨论】:

    标签: apache-kafka


    【解决方案1】:

    使用session.timeout.ms,您只能控制由于心跳而导致的超时,这意味着自上次心跳以来已经过去了session.timeout.ms 毫秒,并且集群将您声明为死节点并触发重新平衡。

    KIP-62 之前,心跳是在轮询中发送的,但现在被移动到特定的后台线程,以避免在您花费超过session.timeout.ms 的时间来调用另一个poll() 时被逐出集群。 将心跳分离到特定线程可以将处理与告诉集群您已启动并正在运行分离,但这会引入进程处于活动状态但没有取得进展的“活锁”情况的风险,因此除了使心跳独立在poll 中引入了一个新的超时,以确保消费者还活着并取得进展。 文档说明了 KIP-62 之前的实现:

    只要消费者在发送心跳,它基本上就会锁定分配给它的分区。如果进程以这样的方式失效,无法取得进展,但仍继续发送心跳,则组中的其他成员将无法接管分区,这会导致延迟增加。然而,心跳和处理都在同一个线程中完成的事实保证了消费者必须取得进展才能保持他们的分配。任何影响处理的停顿也会影响心跳。

    KIP-62 引入的变化包括:

    解耦处理超时:我们建议为记录处理引入一个单独的本地强制超时和一个后台线程以保持会话处于活动状态,直到该超时到期。我们将此新的超时称为“进程超时”,并将其在消费者的配置中公开为 max.poll.interval.ms。此配置设置客户端调用 poll() 之间的最大延迟

    根据您发布的日志,我认为您可能处于这种情况,您的应用程序处理 200 条轮询记录所花费的时间超过 max.poll.interval.ms(默认为 5 分钟)。 如果您在这种情况下,您只能减少更多 max.poll.records 或增加 max.poll.interval.ms

    PD:

    出现在您的日志中的max.poll.interval.ms 配置来自(至少)kafka 0.10.1.0,所以我假设您在那里犯了一个小错误。

    更新

    如果我理解错了,请纠正我,但在您最后的评论中,您说您正在创建 25 个消费者(例如,如果您使用 java,则为 25 个 org.apache.kafka.clients.consumer.KafkaConsumer)并将它们订阅到 N 个不同的主题,但使用相同的 group.id . 如果这是正确的,您将在每次启动或停止KafkaConsumer 时看到重新平衡,因为它将发送包含group.idmember.idJoinGroupLeaveGroup 消息(请参阅相应的kafka protocol)( member.id 不是主机,因此在同一进程中创建的两个消费者仍然具有不同的 ID)。请注意,这些消息不包含主题订阅信息(尽管该信息应该在代理中,但 kafka 不使用它进行重新平衡)。 因此,每次集群收到JoinGroupLeaveGroupgroup.id X 时,都会触发所有具有相同group.id X 的消费者的重新平衡。

    如果您使用相同的 group.id 启动 25 个使用者,您将看到重新平衡,直到创建最后一个使用者并且相应的重新平衡结束(如果您继续看到这种情况,您可能会停止使用者)。

    我有this issue a couple months ago

    如果我们有两个 KafkaConsumer 使用相同的 group.id(在同一个进程或两个不同的进程中运行)并且其中一个已关闭,即使它们订阅了不同的主题,它也会在另一个 KafkaConsumer 中触发重新平衡。 我想经纪人必须只考虑 group.id 进行再平衡,而不是与 LeaveGroupRequest 的对 (group_id,member_id) 对应的订阅主题,但我想知道这是预期的行为还是应该得到改善? 我想这可能是避免在代理中进行更复杂的重新平衡的第一个选择,并且考虑到解决方案非常简单,即只需为订阅不同主题的不同 KafkaConsumer 使用不同的组 ID,即使它们在同一进程中运行。


    当发生重新平衡时,我们会看到重复的消息

    这是预期的行为,一个消费者消费了消息,但在提交偏移量之前触发了重新平衡并且提交失败。当重新平衡完成时,将分配该主题的进程将再次使用该消息(直到提交成功)。

    我被分成两组,现在突然问题在过去 2 小时内消失了。

    您在这里一针见血,但如果您不想看到任何(可避免的)重新平衡,您应该为每个主题使用不同的group.id

    这里是great talk,关于不同的再平衡方案。

    【讨论】:

    • 没有收到您最后的声明you make a little mistake there。具体在哪里?
    • 嗨,明白了。我也对 max.poll.interval.ms 进行了更改。我已将其设置为 600000 毫秒,即 10 分钟。消费者群体仍在重新平衡。现在我不断收到WARN :: auto commit failed
    • 分区分配怎么样?另外,您是否看到至少一条消息正在被消费?
    • 分区分配是自动的。消息被消耗。但是问题太多了。当重新平衡发生时,我们会看到重复的消息。现在,即使在将 max.poll.records 设置为 5 并将 max.poll.interval.ms 设置为 10 分钟之后,我们仍然看到提交失败。然而,我做了一件事。我们从多个主题中消费,并将它们置于单个消费者组中。我分成两组,现在突然问题在过去 2 小时内消失了。知道为什么会发生这种情况吗?我们在该组中有 25 个消费者。全部在单个服务器上
    • 关于您的评论,you will see rebalancing until the last consumer is created。现在每个主题都有不同的CG。但现在我看到,对于某些主题,消费者直到 5 分钟后才开始我的流程重新启动。会不会因为我一个接一个地启动多个消费者而发生?但是还有什么方法可以启动它们呢?
    猜你喜欢
    • 2016-06-08
    • 1970-01-01
    • 2018-12-04
    • 2020-10-28
    • 1970-01-01
    • 2019-07-03
    • 2018-05-05
    • 2021-08-22
    • 1970-01-01
    相关资源
    最近更新 更多