【发布时间】: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