【问题标题】:Kafka KTable - shared aggregation across machinesKafka KTable - 跨机器共享聚合
【发布时间】:2016-08-31 18:20:39
【问题描述】:

假设我有一个包含许多分区的主题。我在其中写入 K/V 数据,并希望在 Tumbling Windows 中按键聚合所述数据。

假设我启动了与分区一样多的工作程序实例,并且每个工作程序实例都在单独的机器上运行。

我将如何确保生成的聚合包含每个键的 all 值? IE 我不希望每个工作实例都有一些值的子集。

这是 StateStore 将用于的东西吗? Kafka 是自己管理这个还是我需要想出一个方法?

【问题讨论】:

    标签: java apache-kafka apache-kafka-streams


    【解决方案1】:

    我将如何确保生成的聚合包含每个键的所有值? IE 我不希望每个工作实例都有一些值的子集。

    一般来说,Kafka Streams 确保相同键的所有值将由相同(且只有一个)流任务处理,这也意味着只有一个应用程序实例(您描述为“工作实例”)将处理该键的值。请注意,一个应用实例可能会运行 1+ 个流任务,但这些任务是隔离的。

    这种行为是通过数据的partitioning实现的,Kafka Streams 确保一个partition 总是由同一个且只有一个流任务处理。键/值的逻辑链接是,在 Kafka 和 Kafka Streams 中,键总是被发送到同一个分区(这里有一个问题,但我不确定是否有必要详细介绍这个问题),因此一个特定的分区 - 在可能的许多分区中 - 包含同一键的所有值。

    在某些情况下,例如在加入两个流 AB 时,您必须确保聚合将在同一个键上操作,以确保来自两个流的数据位于同一个流任务中-- 同样,这一切都是为了确保相关的输入流分区并因此匹配键(分别来自AB)在同一个流任务中可用。您在此处使用的典型方法是selectKey()。一旦完成,Kafka Streams 确保,为了连接两个流 A 和 B 以及创建连接的输出流,相同键的所有值都将由相同的流任务处理,从而由相同的应用程序实例处理。

    例子:

    • A 具有键 userId 和值 { georegion }
    • B 具有键 georegion 和值 { continent, description }

    只有当两个流使用相同的密钥时,才可以加入两个流(从 Kafka 0.10.0 开始)。在此示例中,这意味着您必须对流 A 重新设置密钥(并因此重新分区),以便将生成的密钥从 userId 更改为 georegion。否则,从 Kafka 0.10 开始,您无法加入 AB,因为数据不在负责实际执行联接的流任务中。

    在本例中,您可以通过以下方式对流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 是自己管理这个还是我需要想出一个方法?

    一般来说,不会。上述行为是通过分区实现的,而不是通过状态存储实现的。

    由于您为流定义的操作,有时会涉及到状态存储,这可能解释了您问这个问题的原因。例如,窗口操作将需要管理状态,因此将在幕后创建状态存储。但是您的实际问题——“确保生成的聚合包含每个键的所有值”——与状态存储无关,它与分区行为有关。

    【讨论】:

      【解决方案2】:

      对于工作实例,我假设您是指 Kafka Streams 应用程序实例,对吧? (因为 Kafka Streams 中没有 master/worker 模式——它是一个库而不是一个框架——我们不使用术语“worker”。)

      如果您想按键共同定位数据,则​​需要按键对数据进行分区。因此,当数据从一开始就被写入主题时,您的数据要么由外部生产者按键分区。或者,您在 Kafka Streams 应用程序中显式设置新密钥(例如使用 selectKey()map())并通过调用 through() 重新分配。 (在未来的版本中将不需要显式调用through(),即0.10.1,如果需要,Kafka Streams 将自动重新分发记录。)

      如果消息/记录应该被分区,键不能是null。您还可以通过生产者配置partitioner.class 更改分区架构(请参阅https://kafka.apache.org/documentation.html#producerconfigs)。

      分区完全独立于 StateStore,即使 StateStore 通常用于分区数据。

      【讨论】:

      • 特定键的所有值是否会转到同一个分区?这些都会进入同一个消费者流程吗?我想这就是我要问的。如果每个客户端都聚合每个键的一部分,聚合将不会非常有效。
      • 是的,特定键的所有值都转到同一个分区。是的,这些都将转到同一个流任务,从而转到同一个应用程序实例(您将其描述为“消费者进程”)。正如你所说,如果行为有任何不同,它就不会真正起作用。 ;-) Matthias J. Sax 所描述的是,在某些情况下,例如在加入流时,您可能需要确保两个流具有相同的键(字段),例如通过selectKey()。但是一旦完成,Kafka Streams 会确保同一个键的所有值都转到同一个流任务/应用程序实例。
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-12-20
      • 2020-03-28
      相关资源
      最近更新 更多