【问题标题】:scala.concurrent.Future.onSuccess execution time on different ExecutorServicescala.concurrent.Future.onSuccess 在不同ExecutorService上的执行时间
【发布时间】:2017-05-19 20:27:53
【问题描述】:

我想控制 ExecutionContext 中的线程数。所以我创建了一个 ThreadPoolExecutor 的实例,然后从中创建了 ExecutionContext。

我创建了一些 Futures 并在它们上附加了 onSuccess 回调。我希望每个未来工作完成时都会调用每个 onSuccess 回调。但是我发现所有的 onSuccess 回调都是同时执行的。

import java.util.concurrent.{ Executors, ForkJoinPool }

import scala.concurrent.{ Await, ExecutionContext, Future }
import scala.concurrent.duration.Duration

object Main extends App {
  implicit val ec = ExecutionContext.fromExecutorService(Executors.newFixedThreadPool(2))
  // implicit val ec = ExecutionContext.fromExecutorService(new ForkJoinPool(2))

  val start = System.currentTimeMillis()

  val futures = for {
    i <- 1 to 10
  } yield Future[Int] {
    Thread.sleep(i * 1000)
    i
  }

  futures.foreach { f =>
    f.onSuccess { case i =>
      println(s"${i} Success. ${System.currentTimeMillis() - start}ms elapsed.")
    }
  }

  Await.ready(Future.sequence(futures.toList), Duration.Inf)
  ec.shutdown()
}

// ThreadPoolExecutor Result
// 1 Success. 25060ms elapsed.
// 2 Success. 25064ms elapsed.
// 3 Success. 25064ms elapsed.
// 4 Success. 25064ms elapsed.
// 5 Success. 25064ms elapsed.
// 6 Success. 25064ms elapsed.
// 7 Success. 25065ms elapsed.
// 8 Success. 25065ms elapsed.
// 9 Success. 25065ms elapsed.
// 10 Success. 30063ms elapsed.

// ForkJoinPool Result
// 1 Success. 1039ms elapsed.
// 2 Success. 2036ms elapsed.
// 3 Success. 4047ms elapsed.
// 4 Success. 6041ms elapsed.
// 5 Success. 12042ms elapsed.
// 6 Success. 12043ms elapsed.
// 7 Success. 25060ms elapsed.
// 8 Success. 25060ms elapsed.
// 9 Success. 25060ms elapsed.
// 10 Success. 30050ms elapsed.

上面的结果不是同时打印的。但是当我使用 ForkJoinPool 而不是 ThreadPoolExecutor 时,这个问题得到了缓解。我是否滥用了 ExecutionContext 和 Future?

已编辑:我发现当线程数小于期货数时会出现问题。所以我编辑了上面的代码来重现问题并打印执行时间。

我认为即使线程数很少,也应该按时调用未来的回调......

【问题讨论】:

  • 您需要发布您执行的确切代码。您粘贴的内容不完整,不会产生您描述的输出。事实上,对我来说,一切都如你所愿,而不是你所看到的。
  • 我刚刚编辑了这个问题。当线程数小于期货时,就会出现问题......
  • 您是否希望/需要将每个 Future 标记为 blocking,每个 this StackOverflow post
  • 没有。我希望在未来完成时执行 onSuccess 回调。

标签: scala future executioncontext


【解决方案1】:

我最终知道 Future 回调(onComplete 或 onSuccess)是在提供的 ExecutionContext 的线程上执行的。因此,如果池中没有空闲线程,则无法执行回调。 See scala.concurrent.Future

但我仍然不了解 ForkJoinPool 的行为。我需要研究一下。

【讨论】:

  • ForkJoinPool 默认使用 2*(CPU 逻辑核心)线程。将您的nThreads 替换为相同的数字,您将获得与FixedThreadPool 相同的结果。
猜你喜欢
  • 2011-04-25
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2013-06-29
  • 1970-01-01
  • 2013-06-23
  • 1970-01-01
  • 2017-12-24
相关资源
最近更新 更多