在 Kafka Streams 中,“处理单元”称为 流任务。
任务可以是有状态的和/或无状态的。当发生再平衡事件时,在您的应用程序的一个实例(例如,M)上运行的任务可能会移动到您的应用程序的另一个实例。
主题分区和流任务之间存在 1-1 映射,这保证了一个且只有一个任务将处理来自特定分区的数据。例如,如果任务 3 负责读取和处理分区 P,那么当任务 3 从实例 M 移动到另一个实例 M' 时,M 将停止读取 P(因为它不再运行任务 3),M'(现在运行任务 3)将恢复/继续处理 P。
- 在某个时间 t1,M 收到来自 P 的消息
假设负责处理主题分区P 的流任务称为task(P)。在时间 t1,M 恰好是运行 task(P) 的应用程序实例。这就是上面第 1 点的情况。
- 在稍后的 t2 点,P 从 M 中撤消,因此 M 不再接收来自 P 的消息(例如,因为启动了一个处理 P 的额外工作人员)
在这里,应用程序的另一个实例(您将此实例称为“额外工作人员”)负责运行task(P)。在这里,task(P) 将自动从原始应用实例M 迁移到新实例M'。由task(P) 管理的任何状态(例如,当任务正在执行诸如连接或聚合之类的有状态操作时)当然将与任务一起迁移。在迁移task(P) 时,读取和处理主题分区P 的责任也将从应用实例M 转移到M'。
也许不要想太多“哪个应用实例正在处理主题分区P?”相反,特定分区始终由特定的流任务处理,并且流任务可以跨应用程序实例移动。 (当然,Kafka 的 Streams API 将防止不必要的任务迁移,以确保您的应用程序的处理保持高效。)
- 稍后,t3,M 再次收到来自 P 的消息
这意味着,在时间 t3,M 再次被分配了任务 task(P),因为另一个重新平衡事件 - 可能是因为另一个应用程序实例 M' 被删除,或者发生了其他事情所需的任务迁移。
在 cmets 中被问到这个答案:尽管有一两句话关于州移民也很有用。这不像二进制/物理数据是从一个 RocksDB 实例中获取并传递给另一个。显然,状态是基于容错机制重建的。
有状态的任务使用状态存储来持久化状态信息。这些状态存储是容错的。这些状态存储的真实来源是 Kafka 本身:对状态存储的任何更改(例如,递增计数器)都以流式方式备份到 Kafka——类似于将数据库表的 CDC 流存储到 Kafka 主题中(这些是普通主题,但通常称为“更改日志主题”)。然后,当任务死亡或迁移到另一个容器/VM/机器时,任务的状态存储通过从 Kafka 读取恢复到任务的新容器/VM/机器中(想想:流备份/流恢复) .这会将状态存储恢复到它们在原始容器上的样子,而不会丢失任何数据或重复。
流任务使用 RocksDB 在本地实现状态存储(如在任务的容器中)以进行优化。可以将这些本地 RocksDB 实例视为缓存,就数据安全而言可以丢失,因为如上所述,状态数据的持久存储是 Kafka。