【发布时间】: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