【问题标题】:Akka - Best Approach to configure an Actor that consumes a REST API service (blocking operation)Akka - 配置使用 REST API 服务的 Actor 的最佳方法(阻塞操作)
【发布时间】:2019-10-28 16:36:18
【问题描述】:

我有一个 Akka 消息传递引擎,可以在白天发送数百万条消息,包括 SMS 和电子邮件。我需要引入一种新型消息传递(PushNotification),它包括让每个请求都使用一个 REST API(它还将处理数百万个)。我相信使用 Web 服务是一个阻塞操作,所以从我所读的内容来看,我需要为这个新的 Actor 添加一个单独的调度程序,我的问题是,它是否一定需要是一个带有固定池的线程池执行器-这里提到的大小? (请参阅https://doc.akka.io/docs/akka-http/current/handling-blocking-operations-in-akka-http-routes.html)或者是否可以使用 fork-join-executor 代替?另外,为了不影响当前的两种消息传递方式,最好的方法是什么? (SMS 和 EMAIL)我的意思是如何避免饿死他们的线程池?目前 EMAIL 使用单独的 Dispatcher,而 SMS 使用的是 Default Dispatcher。除了使用阻塞操作(调用 WebService)为 Actor 创建一个新的 Dispatcher 之外,还有其他方法吗?就像创建响应式 Web 服务一样?

【问题讨论】:

  • 您使用的是什么 REST 客户端? Akka Http 使用相同的线程池异步工作,并与 Akka Actors 很好地集成。
  • @Tim,我使用的是 Jersey (com.sun.jersey.api.client.ClientResponse),你的意思是 Akka Http 是响应式的吗?

标签: scala concurrency parallel-processing akka


【解决方案1】:

使用来自 Web 服务的 RESTful API 不必阻塞。

使用来自参与者的 RESTful API 的一种简单方法是使用 Akka HTTP Client。这允许您发送 HTTP 请求并将结果作为消息发送回使用 pipeTo 方法的参与者。

这是一个非常精简的示例(根据文档中的示例稍作修改)。

import akka.http.scaladsl.Http

object RestWorker {
  def props(replyTo: ActorRef): Props =
    Props(new RestWorker(replyTo))
}

class RestWorker(replyTo: ActorRef) extends Actor
{
  implicit val ec: ExecutionContext = context.system.dispatcher

  override def preStart() = {
    Http(context.system).singleRequest(HttpRequest(uri = "https://1.2.3.4/resource"))
      .pipeTo(self)
  }

  def receive = {
    case resp: HttpResponse =>

      val response = ??? // Process response

      replyTo ! response

      self ! PoisonPill
  }
}

【讨论】:

  • 底层概念称为java NIO,或非阻塞IO。检查例如stackoverflow.com/questions/25099640/…
  • 底层概念被称为“异步IO”,由Akka Http来选择如何实现。阻塞/非阻塞是一个不同的概念。
猜你喜欢
  • 1970-01-01
  • 2017-11-30
  • 2021-01-12
  • 2016-04-11
  • 1970-01-01
  • 1970-01-01
  • 2013-03-22
  • 2019-11-21
  • 2015-08-08
相关资源
最近更新 更多