【问题标题】:About Future.firstCompletedOf and Garbage Collect mechanism关于 Future.firstCompletedOf 和垃圾收集机制
【发布时间】:2016-04-05 08:15:01
【问题描述】:

我在实际项目中遇到过这个问题,并通过我的测试代码和分析器证明了这一点。我没有粘贴“tl;dr”代码,而是向您展示一张图片,然后对其进行描述。

简单地说,我使用Future.firstCompletedOf 从 2 个Futures 中得到一个结果,这两个都没有共享的东西,也不关心彼此。尽管如此,这是我要解决的问题,垃圾收集器无法回收第一个 Result 对象,直到两个 Futures 完成

所以我真的很好奇这背后的机制。有人可以从较低的层次解释它,或者提供一些提示让我研究一下。

谢谢!

PS:是不是因为他们共享同一个ExecutionContext

**更新**按要求粘贴测试代码

object Main extends App{
  println("Test start")

  val timeout = 30000

  trait Result {
    val id: Int
    val str = "I'm short"
  }
  class BigObject(val id: Int) extends Result{
    override val str = "really big str"
  }

  def guardian = Future({
    Thread.sleep(timeout)
    new Result { val id = 99999 }
  })

  def worker(i: Int) = Future({
    Thread.sleep(100)
    new BigObject(i)
  })

  for (i <- Range(1, 1000)){
    println("round " + i)
    Thread.sleep(20)
    Future.firstCompletedOf(Seq(
      guardian,
      worker(i)
    )).map( r => println("result" + r.id))
  }

  while (true){
    Thread.sleep(2000)
  }
}

【问题讨论】:

  • 我很好奇你是如何设法证明“结果”不能被垃圾收集的,因为我会说相反,这可能很有趣。也许添加更多关于你如何验证的细节?
  • 显示代码。几乎不可能说出没有它会发生什么。
  • 其实这个问题是普遍性的,不依赖于具体的用例,所以很有可能在没有进一步细节的情况下回答。
  • @GiovanniCaporaletti 和 TheArchetypalPaul 感谢您的回复。稍后我将粘贴我的代码。然而,就像 Regis 所说的那样,非常简单明了,只有 2 个期货和一个分析器。我看到了这个对象并触发了一个 GC 事件,它仍然存在。在我的项目中,“结果”要大得多,而且当发生如此频繁时,它给了我一个 OutOfMemory 错误。

标签: scala concurrency garbage-collection jvm future


【解决方案1】:

让我们看看firstCompletedOf是如何实现的:

def firstCompletedOf[T](futures: TraversableOnce[Future[T]])(implicit executor: ExecutionContext): Future[T] = {
  val p = Promise[T]()
  val completeFirst: Try[T] => Unit = p tryComplete _
  futures foreach { _ onComplete completeFirst }
  p.future
}

在执行{ futures foreach { _ onComplete completeFirst } 时,函数completeFirst 保存在某处 通过ExecutionContext.execute。这个函数到底保存在哪里无关紧要,我们只知道它必须保存在某个地方 以便稍后可以在线程可用时在线程池中选择并执行它。只有当 future 完成时,才不再需要对 completeFirst 的引用。

因为completeFirst 关闭了p,只要还有一个未来(来自futures)等待完成,就有一个对p 的引用可以防止它被垃圾收集(即使通过那一点很可能firstCompletedOf已经返回,从堆栈中删除p

当第一个future 完成时,它会将结果保存到promise 中(通过调用p.tryComplete)。 因为promise p 持有结果,所以至少只要p 可达,结果就可达,并且正如我们所见,只要futures 的至少一个未来尚未完成,p 就可达。 这就是为什么在所有期货都完成之前无法收集结果的原因。

更新: 现在的问题是:它可以修复吗?我认为可以。我们所要做的就是确保第一个未来以线程安全的方式完成对 p 的引用,这可以通过使用 AtomicReference 的示例来完成。像这样的:

def firstCompletedOf[T](futures: TraversableOnce[Future[T]])(implicit executor: ExecutionContext): Future[T] = {
  val p = Promise[T]()
  val pref = new java.util.concurrent.atomic.AtomicReference(p)
  val completeFirst: Try[T] => Unit = { result: Try[T] =>
    val promise = pref.getAndSet(null)
    if (promise != null) {
      promise.tryComplete(result)
    }
  }
  futures foreach { _ onComplete completeFirst }
  p.future
}

我已经对其进行了测试,并且正如预期的那样,它确实允许在第一个未来完成后立即对结果进行垃圾收集。它应该在所有其他方面表现相同。

【讨论】:

  • 感谢您为我完成了这项工作,我盯着firstCompletedOf 看了很长时间,无法弄清楚。而且,结论还是很违背直觉的,不知道有没有人抱怨过……
  • 我添加了一个替代实现来解决这种情况。让我知道它是否适合您(这可能需要向标准库提出拉取请求)。
  • 正如我所观察到的,它工作得很好。线程仍然被占用,但这完全是另一回事。感谢您的帮助!
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2012-03-21
  • 2013-01-15
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多