【问题标题】:KafkaStreams processor partition re-assignment behaviorKafka Streams 处理器分区重新分配行为
【发布时间】:2018-04-28 14:56:50
【问题描述】:

假设我有一个基本的 KafkaStreams 应用程序,其中包含一个主题(具有多个分区)和一种处理消息的处理器类型,如下所示:

    builder.stream(topic)
           .process(() -> new MyProcessor());

以下情况是否会发生?对于 MyProcessor 的特定实例,例如 M(即通过调用处理器供应商获得的特定 java 对象),对于主题的特定分区,例如 P

  1. 在某个时间 t1M 收到来自 P 的消息
  2. 在稍后的时间点 t2PM 被撤销,因此 M 不会收到来自P 不再(例如,因为启动了一个额外的工作人员来处理 P
  3. 稍后,t3M 再次收到来自 P 的消息

我检查了documentation 关于流任务如何与 Kafka 主题分区相关的详细信息,但我没有找到有关这与处理器实例的构造和删除和/或(取消)将主题分区分配给现有发生再平衡时的处理器。

【问题讨论】:

    标签: apache-kafka-streams


    【解决方案1】:

    在 Kafka Streams 中,“处理单元”称为 流任务

    任务可以是有状态的和/或无状态的。当发生再平衡事件时,在您的应用程序的一个实例(例如,M)上运行的任务可能会移动到您的应用程序的另一个实例。

    主题分区和流任务之间存在 1-1 映射,这保证了一个且只有一个任务将处理来自特定分区的数据。例如,如果任务 3 负责读取和处理分区 P,那么当任务 3 从实例 M 移动到另一个实例 M' 时,M 将停止读取 P(因为它不再运行任务 3),M'(现在运行任务 3)将恢复/继续处理 P

    1. 在某个时间 t1,M 收到来自 P 的消息

    假设负责处理主题分区P 的流任务称为task(P)。在时间 t1,M 恰好是运行 task(P) 的应用程序实例。这就是上面第 1 点的情况。

    1. 在稍后的 t2 点,P 从 M 中撤消,因此 M 不再接收来自 P 的消息(例如,因为启动了一个处理 P 的额外工作人员)

    在这里,应用程序的另一个实例(您将此实例称为“额外工作人员”)负责运行task(P)。在这里,task(P) 将自动从原始应用实例M 迁移到新实例M'。由task(P) 管理的任何状态(例如,当任务正在执行诸如连接或聚合之类的有状态操作时)当然将与任务一起迁移。在迁移task(P) 时,读取和处理主题分区P 的责任也将从应用实例M 转移到M'

    也许不要想太多“哪个应用实例正在处理主题分区P?”相反,特定分区始终由特定的流任务处理,并且流任务可以跨应用程序实例移动。 (当然,Kafka 的 Streams API 将防止不必要的任务迁移,以确保您的应用程序的处理保持高效。)

    1. 稍后,t3,M 再次收到来自 P 的消息

    这意味着,在时间 t3,M 再次被分配了任务 task(P),因为另一个重新平衡事件 - 可能是因为另一个应用程序实例 M' 被删除,或者发生了其他事情所需的任务迁移。

    在 cmets 中被问到这个答案:尽管有一两句话关于州移民也很有用。这不像二进制/物理数据是从一个 RocksDB 实例中获取并传递给另一个。显然,状态是基于容错机制重建的。

    有状态的任务使用状态存储来持久化状态信息。这些状态存储是容错的。这些状态存储的真实来源是 Kafka 本身:对状态存储的任何更改(例如,递增计数器)都以流式方式备份到 Kafka——类似于将数据库表的 CDC 流存储到 Kafka 主题中(这些是普通主题,但通常称为“更改日志主题”)。然后,当任务死亡或迁移到另一个容器/VM/机器时,任务的状态存储通过从 Kafka 读取恢复到任务的新容器/VM/机器中(想想:流备份/流恢复) .这会将状态存储恢复到它们在原始容器上的样子,而不会丢失任何数据或重复。

    流任务使用 RocksDB 在本地实现状态存储(如在任务的容器中)以进行优化。可以将这些本地 RocksDB 实例视为缓存,就数据安全而言可以丢失,因为如上所述,状态数据的持久存储是 Kafka。

    【讨论】:

    • 嗨迈克尔,谢谢你的详细解释。我想避免的情况如下:假设处理器实例 M 在流框架之外执行一些操作并为此缓存一些值。如果一个streamtask task(P)首先由M处理,然后在一段时间内不被M处理,则M中的缓存值> 陈旧导致错误行为。将任务(P)重新分配给 M 是否正确将始终触发一个 Processor#close,然后是一个 Processor#init 调用,然后可用于清除 M 的缓存?
    • 是的,没错。当任务不再分配给线程时会调用close(),分配任务时会调用init()。
    • 好的,很好,我认为问题已得到解答,我可以依靠 close/init 使缓存无效。我仍然没有得到一些东西:假设处理器的每个实例只会处理 1 个分区/任务,是否有特殊原因导致处理器实例(java 对象)被关闭并重新初始化刚刚删除和重建?因为随着时间的推移,“工人”似乎可以为不再服务的任务积累大量不活动的处理器实例(最坏的情况#partitions - 1)。
    • 一个应用程序实例可以处理零个、一个或多个流任务(多个分区)。正在从实例中删除任务,但这与处理器中公开的 close() 和 init() 方法是正交的。
    • @MichaelG.Noll Of course, Kafka's Streams API will prevent unnecessary task migrations to ensure the processing of your application stays efficient. 阈值是多少?什么触发了任务迁移?
    猜你喜欢
    • 1970-01-01
    • 2018-10-08
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-03-15
    • 1970-01-01
    相关资源
    最近更新 更多