【问题标题】:Kafka Streams: a map not repartitioning co-partitioned dataKafka Streams:不重新分区共同分区数据的地图
【发布时间】:2018-07-19 11:36:51
【问题描述】:

我有一个来自底层主题的 KStream,其类型为 [K3, V]。 K3是由三个字段组成的密钥,即K3(a,b,c)。然而,主题仅由键字段的子集划分,即 K2 (a,b)。

现在,我想创建一个 KTable 来连接并在我的 PAPI 处理器中使用。我希望这个 KTable 按 K2(a,b) 聚合。聚合只是将值收集到一个集合中。

为此,我必须使用“映射”功能将我的密钥从 K3 转换为 K2。这将(尝试)通过创建新的重新分区主题来重新分区数据(尽管实际上数据将保留在相同的分区中,因为它还将使用 K2 作为分区键),请参阅下面拓扑中的“test-customerStoreName-repartition”。

  Sub-topology: 0
Source: KSTREAM-SOURCE-0000000000 (topics: [test-customerz])
  --> KSTREAM-MAP-0000000003
Processor: KSTREAM-MAP-0000000003 (stores: [])
  --> KSTREAM-FILTER-0000000006
  <-- KSTREAM-SOURCE-0000000000
Processor: KSTREAM-FILTER-0000000006 (stores: [])
  --> KSTREAM-SINK-0000000005
  <-- KSTREAM-MAP-0000000003
Sink: KSTREAM-SINK-0000000005 (topic: test-customerStoreName-repartition)
  <-- KSTREAM-FILTER-0000000006

有没有一种方法可以进行这种聚合,而无需通过地图重新分区?

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    使用 DSL,这是不可能的,因为您无法告诉库不需要重新分区。有一个KIP提议增加这样的功能:https://cwiki.apache.org/confluence/display/KAFKA/KIP-759%3A+Unneeded+repartition+canceling

    您需要直接使用处理器 API,因为处理器 API 没有任何自动重新分区。

    你也可以“破解”一些东西:在map() 之后,返回的KStream 可以转换为KStreamImpl 类型,然后通过反射你可以将内部标志repartitionRequired 设置为false。但这是一个 hack!

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-04-28
      • 1970-01-01
      • 1970-01-01
      • 2018-10-08
      • 1970-01-01
      • 2018-04-16
      • 1970-01-01
      • 2018-10-24
      相关资源
      最近更新 更多