【问题标题】:Kafka streams: Read from ALL partitions in every instance of an applicationKafka 流:从应用程序的每个实例中的所有分区读取
【发布时间】: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


【解决方案1】:

要么将输入主题“data_in”的分区数更改为 1 个分区,要么使用GlobalKtable 从主题中的所有分区获取数据,然后您可以将其加入流。这样,您的应用实例就不必再位于不同的消费者组中了。

代码将如下所示:

private GlobalKTable<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));

  thrashStream.to("new_data_in"); // by sending to an other topic you're forcing a repartition on that topic

  KStream<String, Data> newTrashStream = getBuilder().stream("new_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 = newTrashStream.groupByKey();

  Materialized<String, theDataList, KeyValueStore<Bytes, byte[]>> materialized = Materialized.as("agg-stream-store");
  materialized = materialized.withValueSerde(SerDes.theDataDataListSerde);

// Return a KTable
  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)
  .to("agg_data_in");

  return getBuilder().globalTable("agg_data_in");
}

编辑:我编辑了上面的代码以强制对名为“new_data_in”的主题进行重新分区。

【讨论】:

  • 您确定不会使用 GlobalKTable 覆盖数据吗?
  • 恐怕随着数据进来,当ID相同时GlobalKTable中的旧数据会被覆盖
  • 我的意思是,每个实例只会在每次更新全局表时覆盖其他实例的聚合。
  • 通常当数据不断进入主题“data_in”时,它会计算您创建的聚合方法并将结果写入新主题“agg_data_in”。之后,您必须在您提到的其他主题的流和“agg_data_in”中的 globalktable 之间进行连接。
  • 谢谢!询问来自 agg_data_in 的 GlobalKTable 更新吗?据我所知,GlobalKTable 的更新被覆盖(如果有新数据出现,并在那里找到它的键,它将覆盖旧数据/值)。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-04-27
  • 1970-01-01
  • 1970-01-01
  • 2019-09-19
  • 2020-05-17
相关资源
最近更新 更多