【问题标题】:How to implement Kafka Streams topology that process single topic with interactive queries store and global store如何使用交互式查询存储和全局存储实现处理单个主题的 Kafka Streams 拓扑
【发布时间】:2020-09-15 16:34:28
【问题描述】:

我正在尝试实现 Kafka Streams,它将单个主题流视为具有交互式查询的全局数据库。所以我想拥有:

  1. 记录的全局存储(GlobalKTable、KeyValueStore)

  2. 可查询存储,允许我获取交互式查询的结果(最大)

交互式查询必须计算记录字段之一的全局最大值:

 KStream<String, TercUnitRecord> recordsStream = topologyBuilder.stream(topicName);
 KTable<String, Long> lastUpdateStore = recordsStream.mapValues(record -> record.getLastUpdate())
                .selectKey((key, value) -> "lastdate")
                .groupByKey(Grouped.with(Serdes.String(), Serdes.Long()))
                .reduce((maxValue, currValue) -> maxValue.compareTo(currValue) == 1 ? maxValue : currValue,
 Materialized.as("terc-lastupdate"));

但是,我面临的问题是我无法在一个 Kafka Streams 实例中使用与源相同的单个主题。我已经进行了研究,我发现这样做的唯一方法是通过多个 KafkaStreams 实例,但我不确定这是实现这一目标的正确且唯一的方法。有什么想法吗?

【问题讨论】:

  • .stream(topicName) 使用的是“单个主题”,但在幕后创建了多个主题,如果您是这个意思?
  • 其实是多个店铺,更准确的说。我不想创建新主题,我想拥有内部存储,处理流以获得我想要的(最大)。
  • 您能否指出研究表明您需要多个实例?
  • 如果您的输入主题有多个分区,您将需要使用重新分区主题将所有数据放入一个分区中——否则,您无法计算全局最大值。另请注意,“全局存储”用于“广播状态”——它似乎并不真正适用于您的用例。

标签: java apache-kafka apache-kafka-streams


【解决方案1】:

我为每个任务使用了多个 KafkaStreams 实例,它工作正常。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2019-03-15
    • 2019-09-19
    • 2017-06-09
    • 1970-01-01
    • 1970-01-01
    • 2020-08-28
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多