【发布时间】:2020-08-21 07:46:48
【问题描述】:
// I have hundreds of tasks converting inputs into outputs, which should be persisted.
case class Op(i: Int)
case class Output(i: Int)
val inputs: Seq[Op] = ??? // Number of inputs is huge
def executeLongRunning(op: Op): Output = {
Thread.sleep(Random.nextInt(1000) + 1000) // I cannot predict which tasks will finish first
println("<==", op)
Output(op.i)
}
def executeSingleThreadedSave(outputs: Seq[Output]): Unit = {
synchronized { // Problem is, persisting output is itself a long-running process,
// which cannot be parallelized (internally uses blocking queue).
Thread.sleep(5000) // persist time is independent of outputs.size
println("==>", outputs) // Order of persisted records does not matter
}
}
// TODO: this needs to be implemented
def magicSaver(eventualOutputs: Seq[Future[Output]], saver: Seq[Output] => Unit): Unit = ???
val eventualOutputs: Seq[Future[Output]] = inputs.map((input: Op) => Future(executeLongRunning(input)))
magicSaver(eventualOutputs, executeSingleThreadedSave)
我可以将magicSaver 实现为:
def magicSaver(eventualOutputs: Seq[Future[Output]], saver: Seq[Output] => Unit): Unit = {
saver(Await.result(Future.sequence(eventualOutputs), Duration.Inf))
}
但这有一个主要缺点,即我们在开始持久化输出之前等待所有输入都得到处理,从容错的角度来看,这并不理想。
另一个实现是:
def magicSaver(eventualOutputs: Seq[Future[Output]], saver: Seq[Output] => Unit): Unit = {
eventualOutputs.foreach(_.onSuccess { case output: Output => saver(Seq(output)) })
}
但这会将执行时间延长至 inputs.size * 5secs(由于同步性质,这是不可接受的。
我想要一种方法将已完成的期货组合在一起,当此类期货的数量达到某种权衡大小(例如 100)时,但我不确定如何以简洁的方式完成此操作,无需显式编码轮询逻辑:
def magicSaver(eventualOutputs: Seq[Future[Output]], saver: Seq[Output] => Unit): Unit = {
def waitFor100CompletedFutures(eventualOutputs: Seq[Future[Output]]): (Seq[Output], Seq[Future[Output]]) = {
var completedCount: Int = 0
do {
completedCount = eventualOutputs.count(_.isCompleted)
Thread.sleep(100)
} while ((completedCount < 100) && (completedCount != eventualOutputs.size))
val (completed: Seq[Future[Output]], remaining: Seq[Future[Output]]) = eventualOutputs.partition(_.isCompleted)
(Await.result(Future.sequence(completed), Duration.Inf), remaining)
}
var completed: Seq[Output] = null
var remaining: Seq[Future[Output]] = eventualOutputs
do {
(completed: Seq[Output], remaining: Seq[Future[Output]]) = waitFor100CompletedFutures(remaining)
saver(completed)
} while (remaining.nonEmpty)
}
我在这里缺少任何优雅的解决方案吗?
【问题讨论】:
-
我不确定 stdlib 期货是否提供开箱即用的东西可以优雅地解决这个问题。如果我有这个问题,我会使用 fs2 流。将结果提供给流,在输出可用时保留输出。只是我的 2 美分
-
我认为首先不可能保证准确的结果 - 为了知道完成时间,您必须通过
.mapFuture这可能需要在线程上安排另一个作业水池。因此,您将获得在您的Future完成后触发的Future的完成时间。它可能接近您想要的,但无法保证。