【问题标题】:Is this Akka Kafka Stream configuration benefits from Back Pressure mechanism of the Akka Streams?这个 Akka Kafka Stream 配置是否受益于 Akka Streams 的背压机制?
【发布时间】:2020-12-03 15:26:55
【问题描述】:

我们有一个 Akka 应用程序,它使用 Kafka 主题并将接收到的消息发送到 Akka Actor。我不确定我的编程方式是否使用了 Akka Streams 中内置的背压机制的所有好处。

以下是我的配置...

val control : Consumer.DrainingControl[Done]
Consumer
 .sourceWitOffsetContext(consumerSettings, Subscriptions.topics("myTopic"))
 .map(consumerRecord =>
     val myAvro = consumerRecord.value().asInstanceOf[MyAvro]
     
     val myActor = AkkaSystem.sharding.getMyActor(myAvro.getId)
     
     myActor ! Update(myAvro)          
 )
 .via(Commiter.flowWithOffsetContext(CommitterSettings(AkkaSystem.system.toClassic)))
 .toMat(Sink.ignore)(Consumer.DrainControl.apply)
 .run()

这符合我的业务案例的预期,myActor 收到命令更新(MyAvro)

我对背压的技术概念更恼火,据我所知,背压机制部分由接收器控制,但在这个 Stream 配置中,我的接收器只是“Sink.ignore”。所以我的 Sink 正在为 Back Pressure 做任何事情。

当 Akka Kafka Stream 提交 Kafka Topic 偏移量时,我还很好奇什么?命令发送到 MyActor 邮箱的那一刻?如果是这样,那么我如何处理诸如询问模式之类的场景,Kafka Offset 不应该在询问模式完成之前提交。

我看到一些工厂方法处理手动偏移控制“plainPartitionedManualOffsetSource”、“commitablePartitionManualOffsetSource”,但我找不到任何示例,我可以用我的业务逻辑决定手动提交偏移吗?

作为替代配置,我可以使用类似的东西。

val myActor: ActorRef[MyActor.Command] = AkkaSystem.sharding.getMyActor
val (control, result) =
  Consumer
    .plainSource(consumerSettings, Subscriptions.topics("myTopic"))
    .toMat(Sink.actorRef(AkkaSystem.sharding.getMyActor(myAvro.getId), null))(Keep.both)
    .run()

现在我可以访问 Sink.actorRef,我认为 Back Pressure 机制有机会控制 Back Pressure,自然这段代码不会工作,因为我不知道如何在这个星座下访问 'myAvro'。

谢谢解答..

【问题讨论】:

    标签: akka akka-stream akka-kafka


    【解决方案1】:

    在第一个流中,基本不会有背压。消息发送到myActor 后很快就会发生偏移提交。

    对于背压,您需要等待目标参与者的响应,正如您所说,询问模式是实现这一目标的规范方式。由于从演员外部询问演员(出于所有意图和目的,流在演员之外:阶段由演员执行是一个实现细节)导致Future,这表明mapAsync被要求。

    def askUpdate(m: MyAvro): Future[Response] = ???  // get actorref from cluster sharding, send ask, etc.
    

    然后,您可以将原始流中的 map 替换为

    .mapAsync(parallelism) { consumerRecord =>
      askUpdate(consumeRecord.value().asInstanceOf[MyAvro])
    }
    

    mapAsync 将“飞行中”期货限制为parallelism。如果有 parallelism 期货(当然是由它产生的),它会背压。如果衍生的未来以失败告终(对于请求本身,这通常是超时),它将失败;成功期货的结果(尊重传入订单)将被传递(通常,这些将是akka.Done,尤其是当流中唯一要做的事情是偏移提交和Sink.ignore时)。

    【讨论】:

    • 首先谢谢你的回答,如果我理解正确,如果'parallelism'不大于1,那么将再次没有背压机制。那么下一个问题将是,并行度必须是什么?经典的“核心数 * 2”?哪个最有可能与 Dispatcher 线程池大小相同?以及它如何与 Kafka 的分区交互,这是 Kafka 创建并行性的方式。因此,如果我使用 plainPartitionedManualOffsetSource 而不是 'sourceWitOffsetContext' 并且拥有 4 个 Kafka 分区和并行度 1,我还能获得背压吗?
    • 如果并行度为1,那么从发送请求到收到响应的时间段内仍然会有背压。
    • 分区源和 mapAsync(1) 仍将背压(请注意,消费者仍会从 Kafka 批量读取以提高效率,并且在背压时仍会进行后台轮询以向 Kafka 发出活跃信号)。也就是说,我从来没有真正使用过分区源。
    • “read-from-Kafka-send-asks-to-actors”流中的并行度参数实际上只是控制您背压的阈值。因此,与调度程序池大小一样,它取决于工作负载的具体情况。
    • 例如,您正在请求集群分片 Actor,因此主要的并行性考虑是流的特定分区中的 ID 分布(如果有热点,这通常会降低安全性并行度级别;如果您希望仅在热点中具有背压的高并行度,则可以通过一些自适应限制来改善)。
    【解决方案2】:

    这个说法不正确:

    ...据我所知,背压机制部分由接收器控制,但在此 Stream 配置中,我的接收器只是“Sink.ignore”。所以我的 Sink 正在为 Back Pressure 做任何事情。

    对于背压,Sinks 没有什么特别之处。作为流控制机制的背压将在流中存在异步边界的任何地方自动使用。可能在Sink,但也可能在流中的其他任何地方。

    在您的情况下,您正在连接您的流以与演员交谈。那是你的异步边界,但你这样做的方式是使用 map 并在该地图内使用 ! 与演员交谈。所以没有背压,因为:

    1. map 不是异步操作符,它内部调用的任何内容都不能参与背压机制。所以从 Akka Stream 的角度来看,没有引入异步边界。
    2. ! 是一劳永逸,没有提供关于演员执行任何背压的繁忙程度的反馈。

    就像 Levi 提到的那样,您可以做的是将 tell 更改为 ask 交互,并在其工作完成时让接收参与者响应。然后你可以像 Levi 描述的那样使用mapAsyncmapmapAsync 之间的区别在于 mapAsync 的语义使得它仅在返回 Future 完成时才会向下游发出。即使parallelism 为 1,背压仍然有效。如果您的 Kafka 记录的速度比您的 actor 可以处理的快,mapAsync 将在等待 Future 完成时向上游施加压力。 在这种特殊情况下,我认为增加parallelism 是没有意义的,因为所有这些消息都会被添加到演员的收件箱中,所以这样做不会真正加快任何事情的速度。如果交互是 REST 调用,那么它可以提高整体吞吐量。根据您的参与者处理消息的方式,为mapAsync 增加parallelism 可能会导致吞吐量增加。 paralleslism 值有效地限制了在反压开始之前允许的未完成 Futures 的最大数量。

    【讨论】:

    • 我理解你对 '.map(_ => msg.committableOffset)' 的意思,但我认为(我不确定),'sourceWitOffsetContext' 通过要在 'Commiter 处提交的流传输偏移量.flowWithOffsetContext' 如果 Future 以 Success 结束。
    • 我明白了,这是有道理的。我将编辑答案。还考虑了并行性,尽管我认为如果您的参与者反过来调用某个外部系统,您仍然可以获得一些吞吐量优势,所以当您发送消息时并行性 > 1 并不总是有意义的来自mapAsync 的流中的演员
    • @artur,还要注意集群分片的使用(可能在流中有多个 ID)。在这种情况下,任何参与者邮箱中的消息数量与并行度之间并没有真正有意义的关系,并行度将控制从流中向所有参与者发送的请求数量。
    • 是的@LeviRamsey。谢谢。我错过了。假设它是相同的演员实例。我编辑了答案。
    猜你喜欢
    • 2018-01-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-08-14
    • 2018-11-27
    • 2020-11-05
    • 2018-04-02
    • 1970-01-01
    相关资源
    最近更新 更多