【问题标题】:Scala making parallel network calls using FuturesScala 使用 Futures 进行并行网络调用
【发布时间】:2020-08-30 19:53:14
【问题描述】:

我是 Scala 的新手,我有一个方法,可以从给定的文件列表中读取数据并使用 api 调用 数据,并将响应写入文件。

listOfFiles.map { file =>
  val bufferedSource = Source.fromFile(file)
  val data = bufferedSource.mkString
  bufferedSource.close()
  val response = doApiCall(data)  // time consuming task
  if (response.nonEmpty) writeFile(response, outputLocation)
}

上面的方法,在网络调用过程中耗时太长,所以尝试使用parallel 处理以减少时间。

所以我尝试包装代码块,这会消耗更多时间,但程序很快结束 并且它不会产生任何输出,如上面的代码。

import scala.concurrent.ExecutionContext.Implicits.global

listOfFiles.map { file =>
  val bufferedSource = Source.fromFile(file)
  val data = bufferedSource.mkString
  bufferedSource.close()
  Future {
    val response = doApiCall(data) // time consuming task
    if (response.nonEmpty) writeFile(response, outputLocation)
  }
}

如果您有任何建议,将会很有帮助。 (我也尝试使用“par”,效果很好, 我正在探索“par”以外的其他选项,并使用“akka”、“cats”等框架)

【问题讨论】:

  • 它返回 Either[String, String],我使用模式匹配并提取数据
  • 您的主线程退出,并且由于全局 ExecutionContexts 的线程是守护进程,因此您的进程会在主线程退出时退出。

标签: scala parallel-processing future


【解决方案1】:

基于Jatin,而不是使用包含守护线程的默认执行上下文

import scala.concurrent.ExecutionContext.Implicits.global

使用非守护线程定义执行上下文

implicit val nonDeamonEc = ExecutionContext.fromExecutor(Executors.newCachedThreadPool)

你也可以像这样使用Future.traverse和Await

val resultF = Future.traverse(listOfFiles) { file =>
  val bufferedSource = Source.fromFile(file)
  val data = bufferedSource.mkString
  bufferedSource.close()
  Future {
    val response = doApiCall(data) // time consuming task
    if (response.nonEmpty) writeFile(response, outputLocation)
  }
}

Await.result(resultF, Duration.Inf)

traverse 将List[Future[A]] 转换为Future[List[A]]。

【讨论】:

    猜你喜欢
    • 2020-01-17
    • 1970-01-01
    • 2012-08-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-05-11
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多