【问题标题】:Using reactive-kafka to conditionally process messages使用 reactive-kafka 有条件地处理消息
【发布时间】: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


    【解决方案1】:

    我之前用来解决类似问题的一种方法是利用序列号来保证排序。

    例如,您可以像您描述的那样构建一个流程来保存提交:

    in ~> broadcast ~> isProcessable ~> process ~> producer ~> merge ~> out
       ~> broadcast ~>            isNotProcessable          ~> merge
    

    然后将其包装成这样的订单保留流程(取自我们在我公司开发的库):OrderPreservingFlow。然后可以将生成的流发送到提交者接收器。

    如果您的处理阶段保证排序,您甚至可以通过将逻辑直接嵌入图表来提高效率并避免任何缓冲:

    in ~> injectSeqNr ~> broadcast ~> isProcessable ~> process ~> producer ~> mergeNextSeqNr ~> commit
                      ~> broadcast ~>             isNotProcessable         ~> mergeNextSeqNr
    

    这里你的 mergeNextSeqNr 只是一个修改过的合并阶段,如果输入在端口 1 上可用,如果它的序列号是预期的,你会立即发出它,否则你只需等待数据在另一个端口上可用。

    最终结果应该与使用上面的流程包装完全相同,但如果您嵌入它,您可能更容易适应您的需求。

    【讨论】:

    • 实际上,我刚刚创建了一个prototype 并意识到问题并不存在于 GraphDSL 解决方案中。 broadcast junction actually backpressures 至少有一个 out 分支背压,这意味着在前一个元素离开 Graph 并提交之前没有广播新元素,因此不存在排序问题。
    • 我更新了我的原型并意识到由于生产者是异步的,来自广播结点的背压不起作用。谢谢你的想法。
    猜你喜欢
    • 2018-01-04
    • 2019-08-04
    • 2019-06-21
    • 1970-01-01
    • 2017-04-18
    • 2019-12-15
    • 2022-11-11
    • 2021-01-18
    • 1970-01-01
    相关资源
    最近更新 更多