【问题标题】:Waiting for multiple results in Akka在 Akka 中等待多个结果
【发布时间】:2014-05-11 08:15:36
【问题描述】:
在 Akka 中等待多个参与者结果的正确方法是什么?
Principles of Reactive Programming Coursera 课程有一个使用复制键值存储的练习。在不详细介绍分配的细节的情况下,它需要等待多个参与者的确认,然后才能表明复制已完成。
我使用包含未完成请求的可变映射实现了分配,但我觉得解决方案有“难闻的气味”。我希望有更好的方法来实现看似常见的场景。
为了通过保留我对练习的解决方案来维护课程的荣誉准则,我有一个描述类似问题的抽象用例。
发票行项目需要计算其纳税义务。纳税义务是跨多个税务机关(例如,联邦、州、警察区)应用于该行项目的所有税款的组合。如果每个税务机关都是能够确定项目纳税义务的行为者,则该项目需要所有行为者进行报告,然后才能继续报告整体纳税义务。在 Akka 中完成此场景的最佳/正确方法是什么?
【问题讨论】:
标签:
akka
actor
nonblocking
【解决方案1】:
这是 Akka 中很常见的问题。您有多个演员将为您完成这项工作,您需要将他们组合起来。
Jammie Allen 在他的“Effective Akka”一书中提出的解决方案(它是关于从各种类型的账户中获取银行账户余额)是你产生一个 Actor,它会产生多个能完成这项工作的 Actor(例如计算你的税)。它会等待所有人的回答。
你不应该使用ask,但应该使用tell。
当您生成多个参与者(例如 FederalTaxactor、StateTaxActor...)时,您会向他们发送一条消息,其中包含他们需要处理的数据。然后你知道你需要收集多少个答案。对于每个响应,您都会检查是否所有响应都在那里。如果没有,请稍等。
问题是,如果任何演员失败,您可能会永远等待。因此,您为自己安排了超时消息。如果不是所有答案都存在,则返回操作未成功完成。
Akka 有一个特殊的实用程序,可以作为一个很好的帮手为自己安排超时。
【解决方案2】:
这是我认为您正在寻找的简化示例。它展示了一个像actor这样的master如何产生一些child worker,然后等待他们的所有响应,处理可能发生超时等待结果的情况。该解决方案展示了如何等待初始请求,然后在等待响应时切换到新的接收功能。它还展示了如何将状态传播到等待接收函数中,以避免必须在实例级别具有显式的可变状态。
object TaxCalculator {
sealed trait TaxType
case object StateTax extends TaxType
case object FederalTax extends TaxType
case object PoliceDistrictTax extends TaxType
val AllTaxTypes:Set[TaxType] = Set(StateTax, FederalTax, PoliceDistrictTax)
case class GetTaxAmount(grossEarnings:Double)
case class TaxResult(taxType:TaxType, amount:Double)
case class TotalTaxResult(taxAmount:Double)
case object TaxCalculationTimeout
}
class TaxCalculator extends Actor{
import TaxCalculator._
import context._
import concurrent.duration._
def receive = waitingForRequest
def waitingForRequest:Receive = {
case gta:GetTaxAmount =>
val children = AllTaxTypes map (tt => actorOf(propsFor(tt)))
children foreach (_ ! gta)
setReceiveTimeout(2 seconds)
become(waitingForResponses(sender, AllTaxTypes))
}
def waitingForResponses(respondTo:ActorRef, expectedTypes:Set[TaxType], taxes:Map[TaxType, Double] = Map.empty):Receive = {
case TaxResult(tt, amount) =>
val newTaxes = taxes ++ Map(tt -> amount)
if (newTaxes.keySet == expectedTypes){
respondTo ! TotalTaxResult(newTaxes.values.foldLeft(0.0)(_+_))
context stop self
}
else{
become(waitingForResponses(respondTo, expectedTypes, newTaxes))
}
case ReceiveTimeout =>
respondTo ! TaxCalculationTimeout
context stop self
}
def propsFor(taxType:TaxType) = taxType match{
case StateTax => Props[StateTaxCalculator]
case FederalTax => Props[FederalTaxCalculator]
case PoliceDistrictTax => Props[PoliceDistrictTaxCalculator]
}
}
trait TaxCalculatingActor extends Actor{
import TaxCalculator._
val taxType:TaxType
val percentage:Double
def receive = {
case GetTaxAmount(earnings) =>
val tax = earnings * percentage
sender ! TaxResult(taxType, tax)
}
}
class FederalTaxCalculator extends TaxCalculatingActor{
val taxType = TaxCalculator.FederalTax
val percentage = 0.20
}
class StateTaxCalculator extends TaxCalculatingActor{
val taxType = TaxCalculator.StateTax
val percentage = 0.10
}
class PoliceDistrictTaxCalculator extends TaxCalculatingActor{
val taxType = TaxCalculator.PoliceDistrictTax
val percentage = 0.05
}
然后您可以使用以下代码对此进行测试:
import TaxCalculator._
import akka.pattern.ask
import concurrent.duration._
implicit val timeout = Timeout(5 seconds)
val system = ActorSystem("taxes")
import system._
val cal = system.actorOf(Props[TaxCalculator])
val fut = cal ? GetTaxAmount(1000.00)
fut onComplete{
case util.Success(TotalTaxResult(amount)) =>
println(s"Got tax total of $amount")
case util.Success(TaxCalculationTimeout) =>
println("Got timeout calculating tax")
case util.Failure(ex) =>
println(s"Got exception calculating tax: ${ex.getMessage}")
}
【解决方案4】:
在这种情况下,我对流的体验很好。
我用ActorRefs 开始Source,然后通过mapAsync 将带有ask 的消息发送到ActorRefs 并收集对Seq 的响应。
val f = Source(workers)
.mapAsync(USED_THREAD_COUNT)
(actorRef => (actorRef ? QueryState).mapTo[StateResponse]))
.runWith(Sink.seq)
f onComplete { responses =>
// validate and work with responses
}
希望对你有帮助。