【问题标题】:Can only do 4 concurrent futures as maximum in Scala在Scala中最多只能做4个并发期货
【发布时间】:2017-04-14 15:20:21
【问题描述】:

我认为使用期货可以很容易地让我触发一次代码块,但似乎我一次只能有 4 个期货。

这个限制来自哪里,或者我这样使用它是在滥用 Futures?

import scala.concurrent._
import ExecutionContext.Implicits.global
import scala.util.{Failure, Success}
import java.util.Calendar

object Main extends App{

  val rand = scala.util.Random

  for (x <- 1 to 100) {
    val f = Future {
      //val sleepTime =  rand.nextInt(1000)
      val sleepTime =  2000
      Thread.sleep(sleepTime)

      val today = Calendar.getInstance().getTime()
      println("Future: " + x + " - sleep was: " + sleepTime + " - " + today)
      1;
    }
  }

  Thread.sleep(10000)
}

输出:

Future: 3 - sleep was: 2000 - Mon Aug 31 10:02:44 CEST 2015
Future: 2 - sleep was: 2000 - Mon Aug 31 10:02:44 CEST 2015
Future: 4 - sleep was: 2000 - Mon Aug 31 10:02:44 CEST 2015
Future: 1 - sleep was: 2000 - Mon Aug 31 10:02:44 CEST 2015
Future: 7 - sleep was: 2000 - Mon Aug 31 10:02:46 CEST 2015
Future: 5 - sleep was: 2000 - Mon Aug 31 10:02:46 CEST 2015
Future: 6 - sleep was: 2000 - Mon Aug 31 10:02:46 CEST 2015
Future: 8 - sleep was: 2000 - Mon Aug 31 10:02:46 CEST 2015
Future: 9 - sleep was: 2000 - Mon Aug 31 10:02:48 CEST 2015
Future: 11 - sleep was: 2000 - Mon Aug 31 10:02:48 CEST 2015
Future: 10 - sleep was: 2000 - Mon Aug 31 10:02:48 CEST 2015
Future: 12 - sleep was: 2000 - Mon Aug 31 10:02:48 CEST 2015
Future: 16 - sleep was: 2000 - Mon Aug 31 10:02:50 CEST 2015
Future: 13 - sleep was: 2000 - Mon Aug 31 10:02:50 CEST 2015
Future: 15 - sleep was: 2000 - Mon Aug 31 10:02:50 CEST 2015
Future: 14 - sleep was: 2000 - Mon Aug 31 10:02:50 CEST 2015

我希望它们都同时出现。

为了给出一些上下文,我想我可以使用这个结构并通过一个主循环来扩展它,在这个主循环中,它根据从指数分布中提取的值休眠每个循环,以模拟用户到达/执行查询。每次睡眠后,我想通过将查询发送到程序的驱动程序来执行查询(在本例中为 Spark,驱动程序允许多个线程使用它。)有没有比使用 Futures 更明显的方法?

【问题讨论】:

    标签: scala concurrency


    【解决方案1】:

    当您使用import ExecutionContext.Implicits.global时, 它创建与 CPU 数量相同大小的线程池。

    来自ExecutionContext.scala的来源

    默认的ExecutionContext 实现由工作窃取线程池支持。默认, 线程池使用的目标工作线程数等于 [[https://docs.oracle.com/javase/8/docs/api/java/lang/Runtime.html#availableProcessors-- 可用处理器]]。

    还有一个很好的 StackOverflow 问题:What is the behavior of scala.concurrent.ExecutionContext.Implicits.global?

    由于线程池的默认大小取决于CPU的数量,如果你想使用更大的线程池,你必须写类似

    import scala.concurrent.ExecutionContext
    import java.util.concurrent.Executors
    implicit val ec = ExecutionContext.fromExecutorService(Executors.newWorkStealingPool(8))
    

    在执行Future之前。

    (在你的代码中,你必须把它放在for循环之前。)

    请注意,java 8 中添加了工作窃取池,scala 有自己的 ForkJoinPool 进行工作窃取:scala.concurrent.forkjoin.ForkJoinPool vs java.util.concurrent.ForkJoinPool

    此外,如果您希望每个 Future 有一个线程,您可以编写类似的内容

    implicit val ec = ExecutionContext.fromExecutorService(Executors.newSingleThreadExecutor)
    

    因此,以下代码并行执行 100 个线程

    import scala.concurrent._
    import java.util.concurrent.Executors
    
    object Main extends App{
      for (x <- 1 to 100) {
        implicit val ec = ExecutionContext.fromExecutorService(Executors.newSingleThreadExecutor)
        val f = Future {
          val sleepTime =  2000
          Thread.sleep(sleepTime)
    
          val today = Calendar.getInstance().getTime()
          println("Future: " + x + " - sleep was: " + sleepTime + " - " + today)
          1;
        }
      }
    
      Thread.sleep(10000)
    }
    

    除了工作窃取线程池和单线程执行器,还有一些其他的执行器:http://docs.oracle.com/javase/8/docs/api/java/util/concurrent/Executors.html

    详细阅读文档: http://docs.scala-lang.org/overviews/core/futures.html

    【讨论】:

    • 在使用 newSingleThreadExecutor 的后一种情况下,您仍然会遇到这样一个事实,即一个未来可能会阻塞下一个。它甚至比默认的更糟糕。
    • @hbogert 你把Executors.newSingleThreadExecutor 放在Future 之前吗?我更新了问题,请看一下。
    • 默认情况下,它实际上创建了fork-join pool,而不仅仅是fixed-pool
    • @ymonad 是的,这似乎有效。然而,为每个线程初始化一个新的执行器似乎有点浪费(简单的实验表明没有太大的区别)。但是,您对可用的多个 Executor 的引用使我想到了 CachedThreadPool,这在我的情况下是完美的。
    【解决方案2】:

    使用import scala.concurrent.ExecutionContext.Implicits.global 时的默认池确实具有与您机器上的内核一样多的线程。这对于非阻塞代码(没有同步 io/sleep/...)是理想的,但是当您将它用于阻塞代码时可能会出现问题甚至导致死锁。

    但是,如果您在 scala.concurrent.blocking 块中标记阻塞代码,该池实际上会增长。例如,当您使用在等待 Future 时阻塞的 Await.resultAwait.ready 函数时使用相同的标记。

    请参阅blocking 的 api 文档

    所以你所要做的就是更新你的例子:

    import scala.concurrent.blocking
    ...
    val sleepTime = 2000
    blocking{
      Thread.sleep(sleepTime)
    }
    ...
    

    现在所有期货都将在 2000 毫秒后结束

    【讨论】:

      【解决方案3】:

      你也可以使用

      `implicit val ec = ExecutionContext.fromExecutorService(ExecutorService.newFixedThreadPool(NUMBEROFTHREADSYOUWANT))`
      

      在 NUMBEROFTHREADSYOUWANT 中,您可以指定要启动的线程数。 这将在 Future 之前使用。

      【讨论】:

      • 这并没有解决更深层次的问题。问题是非 CPU 密集型线程不应占用完整的操作系统线程。在您的解决方案中,如果有 100 个非 CPU 密集型任务,则需要 100 个 FixedThreadPool。
      猜你喜欢
      • 2013-04-27
      • 1970-01-01
      • 2012-02-10
      • 2015-08-13
      • 2015-08-25
      • 2020-07-17
      • 2016-07-15
      • 2017-11-18
      • 1970-01-01
      相关资源
      最近更新 更多