【问题标题】:Migrating from standards Akka actors to Akka Streams with back pressure and throttling从标准 Akka Actor 迁移到具有背压和节流的 Akka Streams
【发布时间】:2020-11-05 05:11:33
【问题描述】:

我们实现了以下逻辑来管理针对不同后端的作业: 一个经理 Actor 被启动。这位演员:

  • 加载针对每个后端所需的配置(可变映射后端名称 -> 后端连接器配置);
  • 加载一个 Actor 池 (RoundRobinPool) 以处理每个后端的作业(可变映射后端名称 -> RoundRobinPool Actor Ref)

当 Manager Actor 接收到请求时,它会从消息中检索后端名称并将其转发到相应的 Actor 池以处理作业(假设已注册此后端的配置)。然后将作业请求的结果从参与者返回给原始发送者(我们使用转发的原因)。

这个逻辑运行得很好,但是后端处理工作很慢,我们处于快速发布者、缓慢消费者的典型案例中,当负载增加时这会引发问题。

经过一些研究,Akka Streams 似乎是可行的方法,因为它允许实现背压和节流,这对我们的使用来说是完美的(例如,限制为每秒 5 个请求)。

这个想法是让 Manager Actor 使用相同的路由逻辑,但用 Source.queue 替换 Actor 池。

在注册 Source.queue 时,它​​会像这样执行:

val queue = Source
  .queue[RunBackendRequest](0, OverflowStrategy.backpressure)
  .throttle(5, 1.second)
  .map(r => runBackendRequest(r))
  .toMat(Sink.ignore)(Keep.left)
  .run())

RunBackendRequest的定义在哪里:

case class RunBackendRequest(originalSender: ActorRef, backendConnector: BackendConnector, request: BackendRequest)

而函数runBackendRequest是这样定义的:

private def runBackendRequest(runRequest: RunBackendRequest): Unit = {
    val connector = BackendConnectorFactory.getBackendConnector(configuration.underlying, runRequest.backendConnector.toConfig(), materializer, environment.asJava)
    Future { connector.doSomeWork(runRequest.request) } map { result =>
      runRequest.originalSender ! Success(result)
    } recover {
      case e: Exception => runRequest.originalSender ! Failure(e)
    }
  }
}

当 Manager Actor 接收到消息时,它将根据消息中包含的目标后端的名称将其“提供”到正确的队列。

因此,我有几个问题:

  • 这是在这个特定用例中使用 Akka Stream 的正确方法,还是我们可以用不同的方式更高效地编写它?
  • 可以在 RunBackendRequest 对象中提供原始发送者的 actorRef 以便在 Flow 中回答请求吗?
  • 有没有办法将 Flow 的结果检索到 Future 中,然后 Manager Actor 可以返回请求本身的结果?

Akka Streams 看起来很强大,但显然有一个学习曲线!

【问题讨论】:

    标签: scala akka akka-stream


    【解决方案1】:

    在我看来,Manager Actor 会造成单点故障。也许值得一试:

    • 原始发件人不断敲击 Akka 流 graph 而不是 Manager actor。确保您将ActorRef 传递到下游,以便可以发回回复
    • 在图表内部,使用partition-then-mergeSubstreams 处理针对不同后端连接器的请求。
    • 在图表的最后一步或后端连接器完成后,回复原始发件人。

    总的来说,Colin's article 很好地介绍了如何使用带有分区和合并的 Akka 流来归档您的目标。

    如果您需要更多说明,请告诉我,我可以相应地更新我的答案。

    【讨论】:

    • 嗨易。对 Akka-stream 来说是全新的,是的,如果您能提供说明,我将不胜感激。我只是不熟悉图形和子流。
    • 嗨,我已经用更多材料更新了答案,如果有帮助,请告诉我:)
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-04-04
    • 2021-01-19
    • 1970-01-01
    • 1970-01-01
    • 2021-01-30
    • 1970-01-01
    相关资源
    最近更新 更多