【发布时间】:2020-04-02 14:16:15
【问题描述】:
我正在尝试使用 Kafka Streams(即不是简单的 Kafka 消费者)从重试主题中读取之前无法处理的事件。我希望从重试主题中消费,如果处理仍然失败(例如,如果外部系统关闭),我希望将事件放回重试主题。因此,我不想立即继续消费,而是在消费前等待一段时间,以免系统中充斥着暂时无法处理的消息。
简化,目前的代码是这样做的,我想给它加一个延迟。
fun createTopology(topic: String): Topology {
val streamsBuilder = StreamsBuilder()
streamsBuilder.stream<String, ArchivalData>(topic, Consumed.with(Serdes.String(), ArchivalDataSerde()))
.peek { key, msg -> logger.info("Received event for key $key : $msg") }
.map { key, msg -> enrich(msg) }
.foreach { key, enrichedMsg -> archive(enrichedMsg) }
return streamsBuilder.build()
}
我曾尝试使用 Window Delay 进行设置,但未能成功。我当然可以在 peek 内睡一觉,但这会导致线程挂起,听起来不是一个非常干净的解决方案。
延迟如何工作的确切细节对我的用例来说并不是非常重要。例如,所有这些都可以正常工作:
- 过去
x秒内关于该主题的所有事件都被一次性消耗掉。在开始/结束消费后,流在再次消费前等待x秒 - 每个事件在被放到主题后
x秒处理 - 流使用消息,每个事件之间的延迟为
x秒
如果有人能提供几行 Kotlin 或 Java 代码来完成上述任何操作,我将不胜感激。
【问题讨论】:
标签: java kotlin apache-kafka apache-kafka-streams