【发布时间】: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