【问题标题】:Delay Kafka processor reading from source topic延迟 Kafka 处理器从源主题读取
【发布时间】:2019-10-10 11:04:31
【问题描述】:

我的拓扑包含两个源主题,这些主题由 Kafka Streams 应用程序中的两个不同处理器读取和处理。一个处理器 A 读取其对应的主题并创建一个持久的本地存储,该存储与拓扑中的另一个处理器 B 共享。

我的问题是,我需要在重新启动后以某种方式暂停处理器 B 处理一小段时间,并让处理器 A 有时间从其主题中读取一些事件,在处理器 B 开始处理之前更新其本地存储。

由于两个处理器属于同一个子拓扑,我不能在 init() 中使用 Thread.sleep,因为这会导致整个应用停止。

那么有没有办法让拓扑中的处理器 B 在重新启动应用程序时等待/停止很短的时间,然后再开始从源主题读取并开始处理事件?

【问题讨论】:

    标签: apache-kafka-streams


    【解决方案1】:

    处理顺序基于记录时间戳。因此,如果 A 处理的记录的时间戳小于 B 处理的记录的时间戳,则将首先处理那些“A 记录”。

    明确暂停一侧没有意义,因为它可能违反处理顺序。只需确保您的输入数据带有正确的时间戳,您就不必担心手动暂停。

    【讨论】:

    • 嗨,Matthias,如果我理解正确,使用 CustomTimestampExtractor 并将时间向后移到特定主题源(至少在恢复启动时)会起到作用。谢谢,我会试试看!
    • 只是为了这个答案的完整性,为了能够在同一任务中将一个流优先于另一个流,您需要将 MAX_TASK_IDLE_MS_CONFIG 参数设置为某个值
    • @ypanag 这不是必须,但是是的,它允许获得更严格的保证以权衡一些处理延迟。
    猜你喜欢
    • 1970-01-01
    • 2021-06-11
    • 2016-10-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多