【发布时间】:2017-07-30 11:25:31
【问题描述】:
我启动了几个异步进程,如果需要,它们反过来可以启动更多进程(想想遍历目录结构或类似的东西)。每个进程都会返回一些东西,最后我想等待所有这些都完成并安排一个函数来处理结果集合。
天真的尝试
我的解决方案尝试使用了可变的ListBuffer(我不断添加我生成的期货)和Future.sequence 来安排一些函数在完成此缓冲区中列出的所有这些期货时运行。
我准备了一个最小的例子来说明这个问题:
object FuturesTest extends App {
var queue = ListBuffer[Future[Int]]()
val f1 = Future {
Thread.sleep(1000)
val f3 = Future {
Thread.sleep(2000)
Console.println(s"f3: 1+2=3 sec; queue = $queue")
3
}
queue += f3
Console.println(s"f1: 1 sec; queue = $queue")
1
}
val f2 = Future {
Thread.sleep(2000)
Console.println(s"f2: 2 sec; queue = $queue")
2
}
queue += f1
queue += f2
Console.println(s"starting; queue = $queue")
Future.sequence(queue).foreach(
(all) => Console.println(s"Future.sequence finished with $all")
)
Thread.sleep(5000) // simulates app being alive later
}
它首先调度f1 和f2 期货,然后f3 将在1 秒后以f1 分辨率调度。 f3 本身将在 2 秒内解决。因此,我期望得到的是以下内容:
starting; queue = ListBuffer(Future(<not completed>), Future(<not completed>))
f1: 1 sec; queue = ListBuffer(Future(<not completed>), Future(<not completed>), Future(<not completed>))
f2: 2 sec; queue = ListBuffer(Future(Success(1)), Future(<not completed>), Future(<not completed>))
f3: 1+2=3 sec; queue = ListBuffer(Future(Success(1)), Future(Success(2)), Future(<not completed>))
Future.sequence finished with ListBuffer(1, 2, 3)
但是,我实际上得到了:
starting; queue = ListBuffer(Future(<not completed>), Future(<not completed>))
f1: 1 sec; queue = ListBuffer(Future(<not completed>), Future(<not completed>), Future(<not completed>))
f2: 2 sec; queue = ListBuffer(Future(Success(1)), Future(<not completed>), Future(<not completed>))
Future.sequence finished with ListBuffer(1, 2)
f3: 1+2=3 sec; queue = ListBuffer(Future(Success(1)), Future(Success(2)), Future(<not completed>))
...这很可能是因为我们等待的期货列表在Future.sequence 的初始调用期间是固定的,以后不会更改。
工作,但丑陋的尝试
最终,我用这段代码让它按照我的意愿行事:
waitForSequence(queue, (all: ListBuffer[Int]) => Console.println(s"finished with $all"))
def waitForSequence[T](queue: ListBuffer[Future[T]], act: (ListBuffer[T] => Unit)): Unit = {
val seq = Future.sequence(queue)
seq.onComplete {
case Success(res) =>
if (res.size < queue.size) {
Console.println("... still waiting for tasks")
waitForSequence(queue, act)
} else {
act(res)
}
case Failure(exc) =>
throw exc
}
}
这按预期工作,最终获得所有 3 个期货:
starting; queue = ListBuffer(Future(<not completed>), Future(<not completed>))
f1: 1 sec; queue = ListBuffer(Future(<not completed>), Future(<not completed>), Future(<not completed>))
f2: 2 sec; queue = ListBuffer(Future(Success(1)), Future(<not completed>), Future(<not completed>))
... still waiting for tasks
f3: 1+2=3 sec; queue = ListBuffer(Future(Success(1)), Future(Success(2)), Future(<not completed>))
finished with ListBuffer(1, 2, 3)
但还是很丑。它只是重新启动Future.sequence 等待,如果它看到在完成时队列长于结果数,希望下次完成时情况会更好。当然,这很糟糕,因为它会耗尽堆栈,并且如果此检查将在创建未来和将其附加到队列之间的一个小窗口中触发,则可能容易出错。
是否可以在不使用 Akka 重写所有内容或诉诸 Await.result 的情况下这样做(I can't actually use,因为我的代码是为 Scala.js 编译的)。
【问题讨论】:
标签: multithreading scala future scala.js