【问题标题】:How to wait for future inside Kafka Stream map()?如何在 Kafka Stream map() 中等待未来?
【发布时间】:2022-12-30 11:30:24
【问题描述】:

我正在用 Java 实现 Spring Boot 应用程序,使用 Spring Cloud Stream 和 Kafka Streams 活页夹。

我需要像这样在 KStream map 方法中实现阻塞操作:

public Consumer<KStream<?, ?>> sink() {
    return input -> input
        .mapValues(value -> methodReturningCompletableFuture(value).get())
        .foreach((key, value) -> otherMethod(key, value));
}

completableFuture.get()抛出异常(InterruptedException、ExecutionException)

如何处理这些异常,以便链式方法不被执行和Kafka 消息未被确认?我不能承受消息丢失,将它发送到死信主题不是一种选择。

有没有更好的方法在map() 内部进行阻塞?

【问题讨论】:

    标签: java apache-kafka apache-kafka-streams spring-cloud-stream-binder-kafka


    【解决方案1】:

    您可以尝试 Kafka Streams 中的分支功能来控制链式方法的执行。例如,这是您可以尝试的伪代码。 您可以以此为起点并根据您的特定用例进行调整。

    final Map<String, ? extends KStream<?, String>> branches = 
    input.split()
         .branch(k, v) -> {
            try {
              methodReturningCompletableFuture(value).get();
              return true;
            }
            catch (Exception e) {
              return false;
            }
          }, Branched.as("good-records"))
          .defaultBranch();
    
    final KStream<?, String> kStream = branches.get("good-records");
    
     kStream.foreach((key, value) -> otherMethod(key, value));
    

    这里的想法是,您只会将没有抛出异常的记录发送到命名分支good-records,其他所有内容都将进入默认分支,我们在此伪代码中忽略它。然后,您仅为那些“好”记录调用额外的链接方法(如此foreach 调用所示)。

    这并没有解决抛出异常后不确认消息的问题。这似乎有点挑战。但是,我对那个用例很好奇。当异常发生并且你处理它时,你为什么不想确认消息?如果不使用 DLT,要求似乎有点严格。这里理想的解决方案是您可能想要引入一些重试,一旦重试耗尽,将记录发送到 DLT,这使得 Kafka Streams 消费者确认消息。然后应用程序移动到下一个偏移量。

    假设 methodReturningCompletableFuture() 返回 Future 对象,调用 methodReturningCompletableFuture(value).get() 只是等待达到默认或配置的超时。因此,在KStream map 操作中等待已经是一个很好的方法。我认为没有其他必要让它进一步等待。

    【讨论】:

    • 感谢您的回答,我正在考虑重试,但我不确定重试的合理次数是多少。我的用例是音频转录,Kafka 流包含音频块,我使用 Vosk API 转录它们。虽然错误的原因可能是错误消息,但也可能是运行时错误。我想知道,如果我使用 RetryTemplate,重试消息之后的消息处理是否会停止? KStream中没有手动确认的方法吗?
    • 当重试发生时,当前正在重试的消息之后的消息只需等待,一旦正在重试的流线程被释放,Kafka Streams 就会将消息传递到您的流。 Kafka Streams 以每条记录为基础处理记录。听起来音频块的顺序很重要,如果不是这样,那么您可以增加并发性(流线程数)以同时处理块。
    • 有一些方法可以访问 Kafka Streams 使用的消费者并进行手动确认,但这太低级并且可能涉及反射的使用(对此不确定,但我记得过去某些字段是私有的因此你最终使用一些反射来访问消费者然后调用手动确认)。
    • 确切地说,顺序很重要,所以我想确保退休将停止处理新消息。非常感谢你的帮助!
    猜你喜欢
    • 1970-01-01
    • 2020-08-11
    • 1970-01-01
    • 2018-09-18
    • 2019-08-19
    • 2019-04-20
    • 1970-01-01
    • 2021-04-24
    • 2020-07-01
    相关资源
    最近更新 更多