【问题标题】:Customize window store implementation in KStream - KStream join自定义 KStream 中的窗口存储实现 - KStream join
【发布时间】:2018-05-23 15:42:01
【问题描述】:

我们需要使用一个非常大的窗口执行 Kstream - Kstream 连接,其中左侧的勾号将触发仅与右侧的最新记录的连接,反之亦然。

这不是默认窗口的工作方式,因为 KStreamKStreamJoinProcessor 内部的 window.fetch 返回的 WindowStoreIterator<V> 是一个可以包含多条记录的迭代器。

特别是,我们注意到RockDBWindowStoreretainDuplicates 属性设置为true,我们希望它设置为false。

我们如何自定义 KStream KStream join 的 store 实现?

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    最简单的方法可能是将代码复制到具有新名称的类中并相应地更改逻辑?另一种可能性应该是将两个流转换为KTables 并进行表-表连接(您需要禁用两个输入 KTables 的缓存。

    但请注意,对于您想要的连接类型,很难正确处理乱序数据。

    【讨论】:

    • 我不得不求助于一个大的 JoinWindow 以确保无论到达顺序是什么(左侧或右侧),输出都会打勾。 ktable ktable join 会给我相同的属性吗?
    • KTable 为双方的输入发出连接结果。这篇博文详细解释了 Kafka Streams 中的连接语义:confluent.io/blog/crossing-streams-joins-apache-kafka
    猜你喜欢
    • 2018-02-23
    • 2018-09-18
    • 1970-01-01
    • 2021-09-18
    • 2017-06-02
    • 1970-01-01
    • 2016-12-30
    • 1970-01-01
    • 2020-10-17
    相关资源
    最近更新 更多