4

以下代码对外部 API 进行了一系列调用。此 API不允许每秒超过 3 个查询syncCallToApi是 api 客户端库提供的一个函数,它执行同步请求并返回结果。

callToApi在保持以下语义的同时,Scala 中每秒并发调用不超过 3 次的最佳方式是什么:

val ids = Seq(12980,2932,3441,42334,980,32,4531,7234)
val results: Seq[Item] = ids.map(id => syncCallToApi(id)) 
println(results)
4

1 回答 1

8

您可以使用Akka Throttler链接2

import akka.actor.{Props, Actor, ActorSystem}
import akka.contrib.throttle.Throttler._
import akka.contrib.throttle.TimerBasedThrottler
import akka.util.Timeout
import java.util.concurrent.TimeUnit
import scala.concurrent.duration._
import akka.pattern.ask
import Api._
import scala.concurrent.{Await, ExecutionContext, Future}
import ExecutionContext.Implicits.global

object Main extends App {
  implicit val timeout = Timeout(1 minute)
  val system = ActorSystem("system")
  val throttler = system.actorOf(Props(new TimerBasedThrottler(new Rate(3, Duration(1, TimeUnit.SECONDS)))))
  val worker = system.actorOf(Props(classOf[ExternalApiActor]))
  throttler ! SetTarget(Option(worker))


  val ids = Seq(12980,2932,3441,42334,980,32,4531,7234)
  val listOfFutures: Seq[Future[Item]] = ids.map { id =>
      ask(throttler, id).mapTo[Item]
    }
  val futureList: Future[Seq[Item]] = Future.sequence(listOfFutures)
  val results: Seq[Item] = Await.result(futureList, 1 minute)
  println(results)
  system.shutdown()
}

class ExternalApiActor extends Actor {
  def receive = {
    case id: Int => sender ! syncCallToApi(id)
  }
}

object Api {
  def syncCallToApi(number: Int): Item = {
    println(s"call the API with $number")
    Item(number)
  }
}

case class Item(id: Int)


//build.sbt
scalaVersion := "2.10.1"

libraryDependencies += "com.typesafe.akka" % "akka-contrib_2.10" % "2.1.4"
于 2013-09-09T19:48:15.597 回答