【问题标题】:Kafka stream aggregate卡夫卡流聚合
【发布时间】:2022-07-12 01:42:56
【问题描述】:

我对 kafka 流聚合有疑问。

我想要的是,对于到达输入主题的每个输入数据,我们都会生成一个新版本的输出聚合 KTable,然后将其连接到第二个主题。

实际上,我们没有 1:1...所以我们在加入第二个主题方面做得不够,我们错过了处理。

我确定问题出在聚合上,因为我在一个主题中编写了聚合的输出,我把消费者放在上面:我确实观察到我没有足够的 KTable 版本正在生成。

我们找到了一些改进的设置:通过使用 Kafka 流配置的 COMMIT_INTERVAL_MS_CONFIG 和 CACHE_MAX_BYTES_BUFFERING_CONFIG 参数,我们有更好的处理速度。

使用这些参数是使聚合方法系统地生成聚合 KTable 版本的正确解决方案吗?如果是,应该设置什么值?

提前感谢您的回答。

这里是聚合和连接的代码:

KGroupedStream<String, GenericRecord> groupedEventStream = eventsSource.groupByKey();
KStream<String, String> resultStream =
        groupedEventStream.aggregate(this::initSensorAggregatedRecord, this::updateSensorAggregatedRecord).leftJoin(secondSource,
            this::bindSecondSource).toStream();

这是我们在 kafka 流配置上设置的设置:

props.put(COMMIT_INTERVAL_MS_CONFIG, 0);
props.put(CACHE_MAX_BYTES_BUFFERING_CONFIG, 0);

问候 CG

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    如果您想为每个传入记录获取一个新窗口,您应该使用滑动窗口sliding windows。它们并不完全符合您的要求,但您可以调整窗口时间以使其适合您。

    【讨论】:

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