【问题标题】:Waiting for multiple results in Akka在 Akka 中等待多个结果
【发布时间】:2014-05-11 08:15:36
【问题描述】:

在 Akka 中等待多个参与者结果的正确方法是什么?

Principles of Reactive Programming Coursera 课程有一个使用复制键值存储的练习。在不详细介绍分配的细节的情况下,它需要等待多个参与者的确认,然后才能表明复制已完成。

我使用包含未完成请求的可变映射实现了分配,但我觉得解决方案有“难闻的气味”。我希望有更好的方法来实现看似常见的场景。

为了通过保留我对练习的解决方案来维护课程的荣誉准则,我有一个描述类似问题的抽象用例。

发票行项目需要计算其纳税义务。纳税义务是跨多个税务机关(例如,联邦、州、警察区)应用于该行项目的所有税款的组合。如果每个税务机关都是能够确定项目纳税义务的行为者,则该项目需要所有行为者进行报告,然后才能继续报告整体纳税义务。在 Akka 中完成此场景的最佳/正确方法是什么?

【问题讨论】:

  • 考虑将税务机关表示为一个对象,而不是参与者,并将行项目的纳税义务计算为未来,如danielwestheide.com/blog/2013/01/09/…
  • 您需要为每个订单项使用一个 Actor 吗? Akka docs 建议针对这种情况评估 Future。

标签: 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}")
    }
    

    【讨论】:

      【解决方案3】:

      正如先前的回答所建议的那样,您可能会发现在这种情况下编写期货的能力很有帮助 - 我所知道的期货(和承诺,它们有些相关)的最佳描述在这里:http://docs.scala-lang.org/overviews/core/futures.html

      这可能有助于解释可组合期货如何满足需求,可能比演员更干净,或者与演员结合。

      【讨论】:

        【解决方案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
        }
        

        希望对你有帮助。

        【讨论】:

          猜你喜欢
          • 1970-01-01
          • 2017-08-14
          • 2018-04-05
          • 2013-12-07
          • 2016-11-13
          • 2020-06-18
          • 2021-03-17
          • 2019-10-07
          • 1970-01-01
          相关资源
          最近更新 更多