【发布时间】:2018-12-11 08:00:08
【问题描述】:
当使用 KTable 时,当实例/消费者的数量等于分区的数量时,Kafka 流不允许实例从特定主题的多个分区中读取。我尝试使用 GlobalKTable 实现这一点,这样做的问题是数据将被覆盖,也无法对其应用聚合。
假设我有一个名为“data_in”的主题,有 3 个分区(P1、P2、P3)。当我运行 Kafka 流应用程序的 3 个实例(I1、I2、I3)时,我希望每个实例从“data_in”的所有分区中读取数据。我的意思是 I1 可以从 P1、P2 和 P3 中读取,I2 可以从 P1、P2 和 P3 中读取,I2 等等。
编辑:请记住,生产者可以将两个相似的 ID 发布到“data_in”中的两个不同分区中。所以当运行两个不同的实例时,GlobalKtable 会被覆盖。
请问,如何实现?这是我的代码的一部分
private KTable<String, theDataList> globalStream() {
// KStream of records from data-in topic using String and theDataSerde deserializers
KStream<String, Data> trashStream = getBuilder().stream("data_in",Consumed.with(Serdes.String(), SerDes.theDataSerde));
// Apply an aggregation operation on the original KStream records using an intermediate representation of a KStream (KGroupedStream)
KGroupedStream<String, Data> KGS = trashStream.groupByKey();
Materialized<String, theDataList, KeyValueStore<Bytes, byte[]>> materialized = Materialized.as("agg-stream-store");
materialized = materialized.withValueSerde(SerDes.theDataDataListSerde);
// Return a KTable
return KGS.aggregate(() -> new theDataList(), (key, value, aggregate) -> {
if (!value.getValideData())
aggregate.getList().removeIf((t) -> t.getTimestamp() <= value.getTimestamp());
else
aggregate.getList().add(value);
return aggregate;
}, materialized);
}
【问题讨论】:
-
您应用的每个实例都必须属于不同的消费者组。
-
我们试过了。在每个实例中,我们都创建了两个构建器(B1、B2)。 B1 从所有分区读取数据,因此每个实例都有唯一的 ID。 B2 从另一个主题的一个分区(不是 data_in)读取数据。后来,当我们加入时,我们得到这个错误
Invalid topology: StateStore agg-stream-store is not added yet因为B1和B2有不同的APP_ID -
我们希望将 data_in 中的所有数据(从所有分区中读取)与另一个主题的一个分区进行连接。
标签: java apache-kafka partitioning apache-kafka-streams