我将如何确保生成的聚合包含每个键的所有值? IE 我不希望每个工作实例都有一些值的子集。
一般来说,Kafka Streams 确保相同键的所有值将由相同(且只有一个)流任务处理,这也意味着只有一个应用程序实例(您描述为“工作实例”)将处理该键的值。请注意,一个应用实例可能会运行 1+ 个流任务,但这些任务是隔离的。
这种行为是通过数据的partitioning实现的,Kafka Streams 确保一个partition 总是由同一个且只有一个流任务处理。键/值的逻辑链接是,在 Kafka 和 Kafka Streams 中,键总是被发送到同一个分区(这里有一个问题,但我不确定是否有必要详细介绍这个问题),因此一个特定的分区 - 在可能的许多分区中 - 包含同一键的所有值。
在某些情况下,例如在加入两个流 A 和 B 时,您必须确保聚合将在同一个键上操作,以确保来自两个流的数据位于同一个流任务中-- 同样,这一切都是为了确保相关的输入流分区并因此匹配键(分别来自A 和B)在同一个流任务中可用。您在此处使用的典型方法是selectKey()。一旦完成,Kafka Streams 确保,为了连接两个流 A 和 B 以及创建连接的输出流,相同键的所有值都将由相同的流任务处理,从而由相同的应用程序实例处理。
例子:
- 流
A 具有键 userId 和值 { georegion }。
- 流
B 具有键 georegion 和值 { continent, description }。
只有当两个流使用相同的密钥时,才可以加入两个流(从 Kafka 0.10.0 开始)。在此示例中,这意味着您必须对流 A 重新设置密钥(并因此重新分区),以便将生成的密钥从 userId 更改为 georegion。否则,从 Kafka 0.10 开始,您无法加入 A 和 B,因为数据不在负责实际执行联接的流任务中。
在本例中,您可以通过以下方式对流A 重新加密/重新分区:
// Kafka 0.10.0.x (latest stable release as of Sep 2016)
A.map((userId, georegion) -> KeyValue.pair(georegion, userId)).through("rekeyed-topic")
// Upcoming versions of Kafka (not released yet)
A.map((userId, georegion) -> KeyValue.pair(georegion, userId))
through() 调用仅在 Kafka 0.10.0 中需要实际触发重新分区,更高版本的 Kafka 会自动为您执行这些操作(即将推出的功能已经完成并在 Kafka trunk 中可用)。
StateStore 会用于此目的吗? Kafka 是自己管理这个还是我需要想出一个方法?
一般来说,不会。上述行为是通过分区实现的,而不是通过状态存储实现的。
由于您为流定义的操作,有时会涉及到状态存储,这可能解释了您问这个问题的原因。例如,窗口操作将需要管理状态,因此将在幕后创建状态存储。但是您的实际问题——“确保生成的聚合包含每个键的所有值”——与状态存储无关,它与分区行为有关。