首先,我认为您可能对 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