【发布时间】:2017-11-02 17:56:27
【问题描述】:
我有一个大的计算大致基于以下模式:
def f1(i:Int):Int = ???
def f2(i:Int):Int = ???
def processA(l: List[Int]) =
l.map(i => Future(f1(i)))
def processB(l: List[Int]) = {
val p = processA(l)
p.map(fut => fut.map(f2))
}
def main() = {
val items = List( /* 1k to 10k items here */ )
val results = processB(items)
results.map(_.onComplete ( ... ))
}
我遇到的问题,如果我的理解是正确的,就是处理是广度优先的。 ProcessA 启动数千个 Future,然后 processB 将在 processA 完成后处理的数千个新 Future 入队。 onComplete 回调将很晚才开始触发...
我想把这个深度优先:processA 的几个 Futures 开始,然后 processB 从那里继续,而不是切换到队列中的其他东西。
可以在 vanilla scala 中完成吗?我应该转向一些可以替代 Futures() 和 ThreadPools 的库吗?
编辑:更详细一点。正如答案中所建议的那样,重写为f1 andThen f2,目前是不可行的。实际上,processA and B 正在做很多其他事情(包括副作用)。 processB 依赖于 ProcessA 的事实是私有的。如果暴露会破坏 SoC。
编辑 2:我想我会放宽一点“原版”约束。有人建议 Akka 流会有所帮助。我目前正在查看 scalaz.Task:任何人的意见?
【问题讨论】:
-
你的意思是
p.map(fut => fut.map(f2))中创建的所有Future只会在val p = processA(l)中创建的每个Future完成后启动?我认为情况不一定如此。 -
能不能把
f1、f2和onComplete参数放在同一个同步函数里? -
@CyrilleCorpet No. f1, f2, processA, processB, onComplete 实际上是相当大的代码片段,在应用程序的不同层中。
-
@Jasper-M 不是所有的期货,但这就是我倾向于看到的。有了这个简化的代码、f1() 中的快速任务和 f2() 中更重的任务,您可以轻松地看到在 f2() 执行之前执行了数百个 f1
-
您的意思是在
f2开始执行之前?还是完结?无论如何,我并不是真正的期货专家。我认为这种行为很大程度上取决于您使用的ExecutionContext实现。
标签: scala threadpool future