【问题标题】:Variable number of Futures可变数量的期货
【发布时间】:2017-06-12 07:56:46
【问题描述】:

我必须在 Scala 中对 2D 列表执行一些操作,并且我正在尝试并行化该任务。

目前我有三个 Futures,每个 Futures 取 N 行矩阵并执行必要的计算。是这样写的:

val future1: Future[List[Int]] = Future { makeCalculations(0, 5) }
val future2: Future[List[Int]] = Future { makeCalculations(6, 10) }
val future3: Future[List[Int]] = Future { makeCalculations(11, 15) }

然后我同时使用 for 推导式启动它们,该推导式生成返回值列表。

问题是,我希望通过将 Int 传递给此函数并让它创建确切数量的期货来使其成为动态的。

我尝试使用 for 理解来生成期货,但它们似乎是按顺序启动的,我希望它们同时启动。我找错地方了吗?有没有更好的方法来做到这一点?

【问题讨论】:

标签: multithreading scala asynchronous multidimensional-array future


【解决方案1】:

首先,我认为您可能对 Futures 的安排方式感到困惑:

然后我同时使用 for 推导式启动它们,该推导式生成返回值列表。

事实上,Future 一被创建就被安排执行(apply 方法被调用)。

让我用一个小代码sn-p来说明:

  import scala.concurrent.ExecutionContext.Implicits.global
  import scala.concurrent.duration.Duration
  import scala.concurrent.{Await, Future}

  val start = System.currentTimeMillis()

  val fs = (1 to 20).grouped(2).map { x =>
    Future {
      val ts = System.currentTimeMillis() - start
      Thread.sleep(1000)
      (x.head, x.head + x.length - 1, ts)
    }
  }

  val res = Future.sequence(fs)

  Await.result(res, Duration.Inf).foreach(
    println
  )

在这里我有一个范围(1 to 20),我将其分成相等的部分并从每个部分创建Future。每个未来都包含它的创建时间戳以及原始范围内的开始和结束索引。

另外你可能会注意到期货内部有一个延迟,所以如果它们按顺序执行,我们会看到开始时间的巨大差异;另一方面,如果futures在同一时间并行启动,则启动时间戳几乎相同。

这是我在机器上得到的结果(前两个数字是索引,第三个数字是相对开始时间戳):

  (1,2,7)
  (3,4,8)
  (5,6,8)
  (7,8,8)
  (9,10,9)
  (11,12,9)
  (13,14,9)
  (15,16,9)
  (17,18,1011)
  (19,20,1011)

如您所见,前 8 个期货同时启动,而第 9 个和第 10 个期货延迟了 1 秒。

为什么会这样?因为我使用了scala.concurrent.ExecutionContext.Implicits.global 执行上下文,默认情况下它的并行度等于处理器内核的数量(在我的例子中是 8 个)。


让我们尝试为执行器提供更高的并行度:

  import java.util.concurrent.Executors
  implicit val ex =
    ExecutionContext.fromExecutor(Executors.newFixedThreadPool(512))
  // same code as before

结果:

  (1,2,8)
  (3,4,8)
  (5,6,16)
  (7,8,13)
  (9,10,13)
  (11,12,14)
  (13,14,14)
  (15,16,14)
  (17,18,15)
  (19,20,16)

正如预期的那样,所有期货几乎同时开始。


最后的测试,让我们有一个并行度为 1 的执行器:

  implicit val ex =
    ExecutionContext.fromExecutor(Executors.newFixedThreadPool(1))

结果:

  (1,2,4)
  (3,4,1005)
  (5,6,2008)
  (7,8,3009)
  (9,10,4010)
  (11,12,5014)
  (13,14,6016)
  (15,16,7016)
  (17,18,8017)
  (19,20,9019)

希望这有助于了解 Future 何时安排执行以及如何控制并行性。有任何问题欢迎在 cmets 中提问。


UPD:关于如何使用 Futures 执行矩阵批处理的更清晰示例。

import scala.concurrent.ExecutionContext.Implicits.global
import scala.concurrent.duration.Duration
import scala.concurrent.{Await, Future}


val m = List(
  Array(1,2,3),
  Array(4,5,6),
  Array(7,8,9),
  Array(1,2,3),
  Array(4,5,6)
)

val n = 3 // desired number of futures
val batchSize = Math.ceil(m.size / n.toDouble).toInt

val fs:Iterator[Future[Int]] = m.grouped(batchSize).map {
  rows:List[Array[Int]] =>
    Future  {
      rows.map(_.sum).sum // any logic for the rows;
                          // sum of the elements
                          // as an illustration
    }
}

// convert Iterator[Future[_]] to Future[Iterator[_]]
val res = Future.sequence(fs)

// print results
Await.result(res, Duration.Inf).foreach(
  println
)

结果:

21 //sum of the elements of the first two rows
30 //... third and fourth row
15 //... single last row

【讨论】:

  • 谢谢你,这清除了一些东西。我认为我的问题主要是语法。现在有了我的理解,我可以生成一个 List(f1,f2,f3),如果我无法访问每个变量,我怎么能做同样的事情?
  • @ggfpc 我不完全清楚你的问题是什么以及你目前想要做什么。我试图在我的回答 (Future.sequence((1 to 20).map(_ => Future(???)))) 中为您提供一种如何在循环中动态创建期货的变体。如果这不是您想要的,请详细说明您的问题(您的“二维列表”和您当前的 for 循环是什么样的,您要实现的功能的签名等)。
  • 我会试试你展示的sn-p。我不得不说我也有点困惑。我想要做的是并行化一个对矩阵执行一些计算的函数(例如添加与当前列相邻的列的值)。我的想法是拥有 N 个 Futures,每个 Future 将有 (Size / N) 行的矩阵可以使用。最后,我会有一个结果列表。目前,我的函数接收第一行和最后一行的索引,这是您在我的原始帖子中看到的,并且它有效。剩下要做的就是以某种方式将那些硬编码的 Future 转换为变量。
  • @ggfpc,好的,请检查更新。我试图让我的例子更清楚。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2014-06-17
  • 1970-01-01
  • 2022-11-25
  • 1970-01-01
  • 2016-09-29
  • 1970-01-01
  • 2015-04-11
相关资源
最近更新 更多