【发布时间】:2016-07-12 22:21:06
【问题描述】:
我有一个 Akka 流,我希望流大约每秒向下游发送消息。
我尝试了两种方法来解决这个问题,第一种方法是让流开始处的生产者在 Continue 消息进入此 Actor 时每秒只发送一次消息。
// When receive a Continue message in a ActorPublisher
// do work then...
if (totalDemand > 0) {
import scala.concurrent.duration._
context.system.scheduler.scheduleOnce(1 second, self, Continue)
}
这工作了一小会儿,然后大量的 Continue 消息出现在 ActorPublisher 演员中,我假设(猜测但不确定)来自下游通过背压请求消息,因为下游可以快速消耗但上游没有产生速度快。所以这个方法失败了。
我尝试的另一种方法是通过背压控制,我在流末尾的ActorSubscriber 上使用MaxInFlightRequestStrategy 将消息数量限制为每秒1 条。这很有效,但是一次大约有三个左右的消息进来,而不是一次只有一个。似乎背压控制并没有立即改变消息进入的速率,或者消息已经在流中排队等待处理。
所以问题是,我怎样才能拥有一个每秒只能处理一条消息的 Akka Stream?
我发现MaxInFlightRequestStrategy 是一种有效的方法,但我应该将批量大小设置为 1,它的批量大小默认为 5,这导致了我发现的问题。现在我正在查看提交的答案,这也是解决问题的一种过于复杂的方法。
【问题讨论】:
-
你考虑过使用
Source.tick吗? -
不,让我看看,谢谢。
-
你也可以试试
throttle。
标签: akka rate rate-limiting akka-stream