如何限制Scala中未处理的期货数量?

yur*_*nov 5 parallel-processing scala

如果有办法限制Scala中未处理的期货数量,我无法资助.例如,在以下代码中:

import ExecutionContext.Implicits.global    
for (i <- 1 to N) {
  val f = Future {
    //Some Work with bunch of object creation
  }
}
Run Code Online (Sandbox Code Playgroud)

如果N太大,它最终会抛出OOM.有没有办法限制未处理的Futures ether的数量与队列般的等待或异常?

Ste*_*man 8

因此,最简单的答案是,您可以创建一个ExecutionContext阻止或限制超出特定限制的新任务的执行.看到这篇博文.有关阻塞Java的更加充实的示例ExecutorService,这是一个示例.[如果需要,可以直接使用它,Maven Central上的库就在这里.]这包含了一些非阻塞ExecutorService,您可以使用工厂方法创建java.util.concurrent.Executors.

将Java ExecutorService转换为Scala ExecutionContext只是ExecutionContext.fromExecutorService( executorService ).因此,使用上面链接的库,您可能会有像...

import java.util.concurrent.{ExecutionContext,Executors}
import com.mchange.v3.concurrent.BoundedExecutorService

val executorService = new BoundedExecutorService(
  Executors.newFixedThreadPool( 10 ), // a pool of ten Threads
  100,                                // block new tasks when 100 are in process
  50                                  // restart accepting tasks when the number of in-process tasks falls below 50
 )

implicit val executionContext = ExecutionContext.fromExecutorService( executorService )

// do stuff that creates lots of futures here...
Run Code Online (Sandbox Code Playgroud)

如果你想要一个ExecutorService可以持续整个应用程序的限制,那就没关系了.但是,如果您在代码中的本地化点创建了大量未来,并且您希望ExecutorService在完成它时关闭它.我在Scala [ maven central ]中定义了贷款模式方法,它们都创建了上下文并在我完成后将其关闭.代码最终看起来像......

import com.mchange.sc.v2.concurrent.ExecutionContexts

ExecutionContexts.withBoundedFixedThreadPool( size = 10, blockBound = 100, restartBeneath = 50 ) { implicit executionContext =>
    // do stuff that creates lots of futures here...

    // make sure the Futures have completed before the scope ends!
    // that's important! otherwise, some Futures will never get to run
}
Run Code Online (Sandbox Code Playgroud)

不是使用ExecutorService直接阻塞的,而是通过强制任务调度(Future-creating)Thread执行任务而不是异步运行来使用减慢速度的实例.你做了一个java.util.concurrent.ThreadPoolExecutor使用ThreadPoolExecutor.CallerRunsPolicy.但ThreadPoolExecutor直接构建起来相当复杂.

更新,更性感,更Scala为中心的替代方案是将Akka Streams作为替代Future同时执行"背压"以防止OutOfMemoryErrors.