【发布时间】:2018-07-26 19:08:22
【问题描述】:
我一直在尝试使用 reactive-kafka,但我遇到了条件处理问题,我没有找到令人满意的答案。
基本上,我正在尝试使用一个包含大量消息(每天大约 100 亿条消息)的 kafka 主题,并且根据消息,然后将我的消息的处理版本推送到另一个主题,我正在努力正确地做到这一点。
我的第一次尝试是这样的:
// This is pseudo code.
Source(ProducerSettings(...))
.filter(isProcessable(_))
.map(process(_))
.via(Producer.flow(producerSettings))
.map(_.commitScalaDsl())
.runWith(Sink.ignore)
这种方法的问题在于,我只在阅读我能够处理的消息时才提交,这显然不是很酷,因为如果我必须停止并重新启动我的程序,那么我必须重新阅读一堆无用的消息,而且由于消息太多,我不能那样做。
然后我尝试通过以下方式使用 GraphDSL:
in ~> broadcast ~> isProcessable ~> process ~> producer ~> merge ~> commit
~> broadcast ~> isNotProcessable ~> merge
这个解决方案显然也不好,因为我无法处理的消息会通过图表的第二个分支并在可处理的消息真正推送到它们的目的地之前被提交,这比第一条消息更糟糕,因为它甚至不保证至少一次交付。
有人知道我该如何解决这个问题吗?
【问题讨论】:
标签: scala apache-kafka akka akka-stream akka-kafka