【问题标题】:Purging messages from stream after joining them in Kafka Stream将消息加入 Kafka Stream 后从流中清除消息
【发布时间】:2021-01-11 03:25:59
【问题描述】:

我正在使用 Kafka Streams 通过键连接来自两个不同 Kafka 主题的两种不同类型的消息。我正在使用Sliding time window。此窗口策略保留来自流的信息,其类型数量与消息是否加入某事无关。

在输入流的吞吐量非常高的情况下,Kafka 为执行连接而创建的主题会快速增长,从而消耗大量磁盘空间。

加入后是否有可能从上述主题中清除消息?这样,我将假设一条消息最多与具有相同密钥的另一条消息连接一次。

非常感谢。

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    更新

    从 2.4.0 版开始,您可以通过 StreamJoined 参数配置流-流连接(参见 https://cwiki.apache.org/confluence/display/KAFKA/KIP-479%3A+Add+StreamJoined+config+object+to+Join)。

    您可以通过Stores 工厂类创建WindowedStoreSupplier,并在您传递给join() 方法的StreamJoined 对象上指定供应商。

    原答案

    您可以通过until() 参数减少保留时间:

    stream1.join(stream2, JoinWindows.of(...).until(/*put retention time here*/);
    

    指定的保留时间将用于本地存储以及基础更改日志主题。请注意,如果更改日志主题已经存在,更改 until() 将不会更新主题配置——您需要手动更新主题配置。

    【讨论】:

    • 感谢您的回答。但是,如何使用此功能清除加入的消息?
    • 你不能——你为什么要这样做?连接的语义定义了记录是否可连接,并且不需要“手动”清除记录
    • 好吧,我想你是对的。此时唯一可行的解​​决方案是缩小加入窗口。谢谢。
    • 谢谢Matthias,state 目录中的state store 怎么样?似乎为每 12 小时窗口创建一个目录(id-lkp-this-join-store.1610193600000,id-lkp-this-join-store.1610236800000,id-lkp-this-join-store.1610280000000.. .etc)这些最终会被清除吗?如何减少留存率?我的 joinWindow 配置有 1 小时的宽限期 JoinWindows.of(Duration.ofMillis(1500)).grace(Duration.ofMillis(3597000))
    • 在旧版本中,您可以使用until()。在较新的版本中,您可以将StreamJoined 参数传递给join() 方法,通过在StreamJoined 上指定WindowStoreSupplier(您可以通过Stores 工厂类创建)来设置它。
    【解决方案2】:

    0.11.0.0 在 AdminClient 中引入了一个新的 API deleteRecords 和一个名为 kafka-delete-records 的脚本,可用于删除给定偏移量之前的所有记录。您可以使用它们来清除不再需要的数据。

    详情请见KIP-107。

    【讨论】:

    • 抱歉,您的回复与我的问题有什么关系?我的意思是,你能给我一个具体的例子吗?
    • 您应该确定中间连接主题和之前没有未处理记录的偏移量,然后将类似 {"partitions":[{"topic": "KSTREAM-JOIN-0000000000", "partition": 0,"offset": <your_offset>}],"version":1} 的 json 文件提供给 kafka-delete-records 脚本。
    • 嗯,在企业应用程序的上下文中有些复杂。
    猜你喜欢
    • 2019-09-14
    • 1970-01-01
    • 1970-01-01
    • 2020-04-30
    • 2015-04-19
    • 2021-04-17
    • 2021-06-12
    • 2023-02-23
    • 2015-01-25
    相关资源
    最近更新 更多