【问题标题】:How to run more than 1 application instances of ktable-ktable joins kafka streams application on single partitioned kafka topics?如何在单个分区的 kafka 主题上运行 1 个以上的 ktable-ktable 加入 kafka 流应用程序的应用程序实例?
【发布时间】:2019-01-27 08:04:46
【问题描述】:
KTable<Key1, GenericRecord> primaryTable = createKTable(key1, kstream, statestore-name);
KTable<Key2, GenericRecord> childTable1 = createKTable(key1, kstream, statestore-name);
KTable<Key3, GenericRecord> childTable2 = createKTable(key1, kstream, statestore-name);

primaryTable.leftJoin(childTable1, (primary, choild1) -> compositeObject)
            .leftJoin(childTable2,(compositeObject, child2) -> compositeObject, Materialized.as("compositeobject-statestore"))
.toStream().to(""composite-topics)

对于我的应用程序,我使用的是 KTable-Ktable 连接,因此每当在主流或子流上接收到数据时,它都可以使用所有三个表的 setter 和 getter 设置它的复合对象。这三个传入的流具有不同的键,但是在创建 KTable 时,我将所有三个 KTable 的键设置为相同。

我有一个分区的所有主题。当我在单个实例上运行应用程序时,一切运行良好。我可以看到compositeObject 填充了来自所有三个表的数据。 通过记录 ID 和本地 statestore 名称,所有交互式查询也可以正常运行。

但是当我运行同一个应用程序的两个实例时,我看到 CompositeObject 包含主要和 child1 数据,但 child2 仍然为空。即使我尝试使用交互式查询调用 statestore,它也不会返回任何内容。

我正在使用 spring-cloud-stream-kafka-streams 库来编写代码。

请提出未设置的原因以及处理此问题的正确解决方案。

【问题讨论】:

    标签: apache-kafka-streams spring-cloud-stream spring-kafka


    【解决方案1】:

    Kafka Streams 的扩展模型与输入主题分区的数量相关。因此,如果您的输入主题是单分区的,则您无法横向扩展。输入主题分区的数量决定了您的最大并行度。

    因此,您需要创建具有更高并行度的新主题。

    【讨论】:

    • 谢谢@Matt。我有 4 个表 - 3 个使用连接,而另一个仅用于查找。所以我现在更改了代码 -
    猜你喜欢
    • 2022-01-09
    • 1970-01-01
    • 2020-11-04
    • 2021-06-08
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多