【问题标题】:Delaying Kafka Streams consuming延迟 Kafka Streams 消费
【发布时间】: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 内睡一觉,但这会导致线程挂起,听起来不是一个非常干净的解决方案。

延迟如何工作的确切细节对我的用例来说并不是非常重要。例如,所有这些都可以正常工作:

  1. 过去x 秒内关于该主题的所有事件都被一次性消耗掉。在开始/结束消费后,流在再次消费前等待 x
  2. 每个事件在被放到主题后x 秒处理
  3. 流使用消息,每个事件之间的延迟为 x

如果有人能提供几行 Kotlin 或 Java 代码来完成上述任何操作,我将不胜感激。

【问题讨论】:

    标签: java kotlin apache-kafka apache-kafka-streams


    【解决方案1】:

    您不能真正暂停使用 Kafka Streams 从输入主题中读取——“延迟”的唯一方法是调用“睡眠”,但正如您所提到的,这会阻塞整个线程并且不是一个好的解决方案。

    但是,您可以做的是使用有状态处理器,例如,process()(带有附加的状态存储)而不是foreach()。如果重试失败,您不会将记录放回输入主题,而是将其放入存储中并注册一个带有所需重试延迟的标点符号。如果标点触发,则重试,如果重试成功,则从存储中删除条目并取消标点;否则,请等到标点符号再次触发。

    【讨论】:

    • 谢谢!但是,我有几个map 步骤(尽管在我的示例中简化为一个),所有这些都可能失败。每个map 会被process 代替吗?
    • 您可以做不同的事情:或者,将所有内容折叠到单个 process() 调用中(似乎是最简单的解决方案),或者从将记录添加到存储中的 transform() 开始,并且如果中间的所有步骤都成功,则有一个最终的transform()(或process())从存储中删除一条记录(注意,您可以通过将其传递给transforms()来共享存储)。
    • @MatthiasJ.Sax 请解释你的意思是注册一个标点符号和所需的重试延迟
    • 不确定您的确切问题是什么。如果您不知道标点符号是什么,请查看文档:docs.confluent.io/platform/current/streams/developer-guide/…
    猜你喜欢
    • 2018-06-08
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-08-29
    • 1970-01-01
    • 2019-06-21
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多