【问题标题】:Kafka Stream Rebalancing : State transition from REBALANCING to ERRORKafka Stream Rebalancing:从 REBALANCING 到 ERROR 的状态转换
【发布时间】:2017-12-28 03:18:42
【问题描述】:

我有 4 个带有单个分区的主题和三个应用程序实例。我试图通过编写一个自定义的 PartitionGrouper 来实现可扩展性,它将创建如下 3 个任务:

第一个实例-topic1,partition0,topic4,partition0

第二个实例-topic2,partition0

第三个实例-topic3,partition0

我将 NUM_STANDBY_REPLICAS_CONFIG 配置为 1,因为它会在本地维护状态(也可以消除 invalidstatestore 异常)。

上述设置适用于两个实例。当我将它增加到三个实例时,我开始面临重新平衡的问题。

StickyTaskAssignor:58 - Unable to assign 1 of 1 standby tasks for task [1009710637_0]. There is not enough available capacity. You should increase the number of threads and/or application instances to maintain the requested number of standby replicas.
    [INFO ] 2017-12-25 20:05:42.221 [app-03-cfaf7841-dc19-4ee4-9d05-ae4928c21449-StreamThread-1] StreamThread:888 - stream-thread [app-03-cfaf7841-dc19-4ee4-9d05-ae4928c21449-StreamThread-1] State transition from PARTITIONS_REVOKED to PARTITIONS_ASSIGNED.
    [INFO ] 2017-12-25 20:05:42.221 [app-03-cfaf7841-dc19-4ee4-9d05-ae4928c21449-StreamThread-1] KafkaStreams:268 - stream-client [app-03-cfaf7841-dc19-4ee4-9d05-ae4928c21449] State transition from REBALANCING to REBALANCING.
    [INFO ] 2017-12-25 20:05:42.276 [app-03-cfaf7841-dc19-4ee4-9d05-ae4928c21449-StreamThread-1] StreamThread:195 - stream-thread [app-03-cfaf7841-dc19-4ee4-9d05-ae4928c21449-StreamThread-1] partition assignment took 55 ms.
    current active tasks: [1009710637_0]
    current standby tasks: [1240464215_0, 1833680710_0]
    previous active tasks: []
    [INFO ] 2017-12-25 20:05:42.631 [app-03-cfaf7841-dc19-4ee4-9d05-ae4928c21449-StreamThread-1] StreamThread:939 - stream-thread [app-03-cfaf7841-dc19-4ee4-9d05-ae4928c21449-StreamThread-1] Shutting down
    [INFO ] 2017-12-25 20:05:42.631 [app-03-cfaf7841-dc19-4ee4-9d05-ae4928c21449-StreamThread-1] StreamThread:888 - stream-thread [app-03-cfaf7841-dc19-4ee4-9d05-ae4928c21449-StreamThread-1] State transition from PARTITIONS_ASSIGNED to PENDING_SHUTDOWN.
    [INFO ] 2017-12-25 20:05:42.633 [app-03-cfaf7841-dc19-4ee4-9d05-ae4928c21449-StreamThread-1] KafkaProducer:972 - Closing the Kafka producer with timeoutMillis = 9223372036854775807 ms.
    [INFO ] 2017-12-25 20:05:42.638 [app-03-cfaf7841-dc19-4ee4-9d05-ae4928c21449-StreamThread-1] StreamThread:972 - stream-thread [app-03-cfaf7841-dc19-4ee4-9d05-ae4928c21449-StreamThread-1] Stream thread shutdown complete
    [INFO ] 2017-12-25 20:05:42.638 [app-03-cfaf7841-dc19-4ee4-9d05-ae4928c21449-StreamThread-1] StreamThread:888 - stream-thread [app-03-cfaf7841-dc19-4ee4-9d05-ae4928c21449-StreamThread-1] State transition from PENDING_SHUTDOWN to DEAD.
    [WARN ] 2017-12-25 20:05:42.638 [app-03-cfaf7841-dc19-4ee4-9d05-ae4928c21449-StreamThread-1] KafkaStreams:343 - stream-client [app-03-cfaf7841-dc19-4ee4-9d05-ae4928c21449] All stream threads have died. The Kafka Streams instance will be in an error state and should be closed.
    [INFO ] 2017-12-25 20:05:42.638 [app-03-cfaf7841-dc19-4ee4-9d05-ae4928c21449-StreamThread-1] KafkaStreams:268 - stream-client [app-03-cfaf7841-dc19-4ee4-9d05-ae4928c21449] State transition from REBALANCING to ERROR.

【问题讨论】:

    标签: apache-kafka-streams


    【解决方案1】:

    我假设您的 PartitionGrouper 破坏了某些东西。编写正确的自定义分区分组器非常困难,因为您需要了解很多有关 Kafka Streams 的内部知识。因此,首先不建议这样做。

    错误本身意味着,StandbyTask 无法成功分配给线程,因为没有足够的线程。一般来说,这个想法是,不能将 StandbyTask 分配给运行相应“活动”任务的线程或相同 StandbyTasks 的另一个副本:它不会增加容错性,只会浪费内存,就好像线程死了一样,所有任务终止。

    为什么你会得到这个错误还不清楚(愉快的调试:))。

    但是,对于您的用例,您应该只启动订阅各个主题的不同应用程序实例并使用不同的application.id 来扩展您的应用程序。

    【讨论】:

    • 谢谢 :) 您建议的另一种方法是将流通过具有许多分区的虚拟主题传递,以根据stackoverflow.com/questions/47325678/… 实现可伸缩性。如果我的自定义流分区器根据字符串的哈希码值将流拆分为多个分区,那么在遇到虚拟主题时,与单个分区相关的输入流是否可能会丢失顺序?
    • 例如:我的输入流有 5 个事件属于 hashcode-107,顺序为:1,3,4,5,2 通过虚拟主题时是否有机会提交将虚拟主题设为 1,2,3,4,5?
    • Depends: order is reserved by key -- 因此,如果单个主题的所有记录都具有相同的键,是的,顺序将被保留(尽管它可能与来自其他主题的记录交错) .如果您没有任何键,则可以在写入虚拟主题之前将主题名称设置为代理键,以确保保留顺序。
    • "一般来说,StandbyTask 不能分配给运行相应“活动”任务或相同 StandbyTasks 的另一个副本的线程:它不会增加容错性,而只会增加浪费内存,就好像一个线程死了,所有的任务都死了。”你能深入了解一下吗?
    • 如果单个线程(或实际上是实例)将维护活动和相应的备用任务,并且实例死亡,您将同时失去活动和备用。因此,备用将不会提供增加的容错/可用性。如果分配它当然仍然会使用资源,但没有任何收益。因此,如果没有能力避免不必要地浪费资源,Kafka Streams 不会分配备用资源,而不是浪费资源。
    猜你喜欢
    • 1970-01-01
    • 2021-06-04
    • 2018-04-20
    • 2019-09-08
    • 2017-05-22
    • 2020-05-08
    • 1970-01-01
    • 1970-01-01
    • 2018-09-23
    相关资源
    最近更新 更多