【问题标题】:Can I avoid repartition in the below kafka stream我可以避免在下面的 kafka 流中重新分区吗
【发布时间】:2022-02-07 19:58:59
【问题描述】:

我正在尝试通过更改键来通过主题流创建一个表,但值保持不变。是否可以避免重新分区?

streamsBuilder.stream(
            TOPIC,
            Consumed.with(IdSerde(), ValueSerde())
        )
            .peek { key, value -> logger.info("Consumed $TOPIC, key: $key, value: $value") }
            .filter { _, value -> value != null }
            .selectKey(
                { _, value -> NewKey(value.newKey.toString()) },
                Named.`as`("changeKey")
            )
            .toTable(
                Materialized.`as`<NewKey, Value, KeyValueStore<Bytes, ByteArray>>(
                    NEW_TABLE_NAME
                )
                    .withKeySerde(NewKeySerde())
                    .withValueSerde(ValueSerde())
            )
        return streamsBuilder.build()

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    我能想到的唯一方法是您的生产者应该使用适当的键将数据发送到输入主题,而不是您在流应用程序中选择键。

    否则,答案是

    在进入解释之前,先了解一下KTable的背景

    KTable 是变更日志流的抽象。这意味着一个 KTable 仅保存给定键的最新值。

    考虑以下四个记录(按顺序)发送到流

    ("alice", 1) --> ("bob", 100) --> ("alice", 3) --> ("bob", 50)
    

    对于上述内容,KTable 如下所示(alice 和 bob 是消息的键):

    ("alice", 3)
    ("bob", 50)
    

    说明

    每当我们运行 kafka 流应用程序的多个实例时, 每个实例中的 KTable 仅保存 应用。因此,需要重新分区以确保 具有相同 key 的消息位于同一分区中。

    为了更好地理解这一点,让我们考虑一个例子。我们有一个包含两个分区的主题。 Partition 0有一个事件,Partition 1有三个事件如下(null代表没有key):

    Topic Partition 0:  (null, {"name": alice,"count": 1}) , 
    Topic Partition 1:  (null, {"name": alice,"count": 3}) , (null, {"name": "bob","count": 100}), (null, {"name": "bob","count": 50})
    

    我们创建一个 kafka 流应用程序来从该主题中读取数据,并使用名称字段作为键创建一个 KTable。此外,流应用程序使用两个实例运行。每个实例将被分配一个分区,如下所示:

    Topic Partition 0: -----> Instance 1
    Topic Partition 1: -----> Instance 2
    

    由于每个实例都在本地维护 KTable,因此如果不进行重新分区 - KTable 将处于两个实例的不一致状态。 如下所示:

    Instance 1  KTable : ("alice", {"count":1})
    Instance 2 KTable : ("alice", {"count":1}), ("bob", {"count":50})
    

    因此,如果 KTable 是在 selectKey 操作之后创建的,为了避免上述 kafka 流重新分区主题。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2020-04-10
      • 2016-11-09
      • 1970-01-01
      • 2021-11-26
      • 2021-09-21
      • 2010-10-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多