对于Kafka Streams,如果我们使用较低级别的处理器API,我们可以控制是否提交.因此,如果我们的代码中出现问题,并且我们不想提交此消息.在这种情况下,Kafka将多次重新发送此消息,直到问题得到解决.
但是如何控制在使用其更高级别的流DSL API时是否提交消息?
资源:
http://docs.confluent.io/2.1.0-alpha1/streams/developer-guide.html
有一个 kafka 集群,我从中消费两个主题并加入它。使用 join 的结果,我对数据库进行了一些操作。对 DB 的所有操作都是异步的,因此它们返回给我一个 Future(scala.concurrent.Future,但无论如何它与 java.util.concurrent.CompletableFuture 相同)。所以结果我得到了这样的代码:
val firstSource: KTable[String, Obj]
val secondSource: KTable[String, Obj2]
def enrich(data: ObjAndObj2): Future[EnrichedObj]
def saveResultToStorage(enrichedData: Future[EnrichedObj]): Future[Unit]
firstSource.leftJoin(secondSource, joinFunc)
.mapValues(enrich)
.foreach(saveResultToStorage)
Run Code Online (Sandbox Code Playgroud)
我可以在流中使用未来值进行操作,还是有更好的方法来处理异步任务(例如 Akka 流中的 .mapAsync)?