【发布时间】:2019-10-10 11:04:31
【问题描述】:
我的拓扑包含两个源主题,这些主题由 Kafka Streams 应用程序中的两个不同处理器读取和处理。一个处理器 A 读取其对应的主题并创建一个持久的本地存储,该存储与拓扑中的另一个处理器 B 共享。
我的问题是,我需要在重新启动后以某种方式暂停处理器 B 处理一小段时间,并让处理器 A 有时间从其主题中读取一些事件,在处理器 B 开始处理之前更新其本地存储。
由于两个处理器属于同一个子拓扑,我不能在 init() 中使用 Thread.sleep,因为这会导致整个应用停止。
那么有没有办法让拓扑中的处理器 B 在重新启动应用程序时等待/停止很短的时间,然后再开始从源主题读取并开始处理事件?
【问题讨论】: