【问题标题】:scala Future processing depth-first not breadth-firstscala 未来处理深度优先而不是广度优先
【发布时间】: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完成后启动?我认为情况不一定如此。
  • 能不能把f1f2onComplete参数放在同一个同步函数里?
  • @CyrilleCorpet No. f1, f2, processA, processB, onComplete 实际上是相当大的代码片段,在应用程序的不同层中。
  • @Jasper-M 不是所有的期货,但这就是我倾向于看到的。有了这个简化的代码、f1() 中的快速任务和 f2() 中更重的任务,您可以轻松地看到在 f2() 执行之前执行了数百个 f1
  • 您的意思是在f2 开始执行之前?还是完结?无论如何,我并不是真正的期货专家。我认为这种行为很大程度上取决于您使用的 ExecutionContext 实现。

标签: scala threadpool future


【解决方案1】:

我不是 100% 确定我理解了这个问题,因为 processB (f2) 在 processA (f1) 的结果之上运行,您不能在尚未由 f1 计算的值上调用 f2,所以我的回答是基于以下假设:

  • 您想限制进行中的工作
  • 您想在f1 之后立即执行f2

所以这里有一个解决方案:

import scala.concurrent._
def process(noAtATime: Int, l: List[Int])(transform: Int => Int)(implicit ec: ExecutionContext): Future[List[Int]] = {
  // define an inner async "loop" to process one chunk of numbers at a time
  def batched(i: Future[Iterator[List[Int]]], result: List[List[Int]]): Future[List[Int]] =
    i flatMap { it =>
      // if there are more chunks to process
      // we process all numbers in the chunk as parallel as possible,
      // then combine the results into a List again, then when all are done,
      // we recurse via flatMap+batched with the iterator
      // when we have no chunks left, then we un-chunk the results
      // reassemble it into the original order and return the result
      if(it.hasNext) Future.traverse(it.next)(n => Future(transform(n))).flatMap(re => batched(i, re :: result))
      else Future.successful(result.reverse.flatten) // Optimize this as needed
    }
  // Start the async "loop" over chunks of input and with an empty result
  batched(Future.successful(l.grouped(noAtATime)), List.empty)
}


scala> def f1(i: Int) = i * 2 // Dummy impl to prove it works
f1: (i: Int)Int

scala> def f2(i: Int) = i + 1 // Dummy impl to prove it works
f2: (i: Int)Int

scala> process(noAtATime = 100, (1 to 10000).toList)(n => f2(f1(n)))(ExecutionContext.global)
res0: scala.concurrent.Future[List[Int]] = Future(<not completed>)

scala> res0.foreach(println)(ExecutionContext.global)

scala> List(3, 5, 7, 9, 11, 13, 15, 17, 19, 21, 23, 25, 27, 29, 31, 33, 35, 37, 39, 41, 43, 45, 47, 49, 51, 53, 55, 57, 59, 61, 63, 65, 67, 69, 71, 73, 75, 77, 79, 81, 83, 85, 87, 89, 91, 93, 95, 97, 99, 101, 103, 105, 107, 109, 111, 113, 115, 117, 119 …

如果您愿意并且能够使用更适合当前问题的库,have a look at this reply

【讨论】:

  • 这会等待整个块完成,然后再开始下一个块。对于比在 int 上进行基本数学运算稍微复杂一点的事情,它的效果并不好。
  • 我不明白你怎么能得出这样的结论——因为它可以分块和取消分块,所以它绝对可以任意嵌套地做这个——任意多的块。
  • “分块和不分块”?不知道什么意思,不好意思。你可以随意嵌套,但事实是在前一个块完成之前,下一个块不会开始。所以,如果你正在运行,比如说,16 个核心,块大小为 16,15 个任务在 10 秒内完成,一个需要一分钟,那么你有 15 个核心空闲 50 秒
  • 您可以决定要同时运行多少个,这就是它的意思。 IE。在我的示例中,它只在一个块之后运行,但这不是固有的限制,因为代码是完全可嵌套的。
  • 固有的限制是没有办法让N个任务在任何时候都被执行。这与您并行运行多少块无关,而是关于整个块的完成是开始另一个块的条件。
【解决方案2】:

您的问题最好用流表示。作业进入流并被处理,背压用于确保一次只完成有限数量的工作。在 Akka 流中,它看起来像这样:

Source(items)
  .mapAsync(4)(f1)
  .mapAsync(4)(f2)
  .<whatever you want to do with the result>

需要仔细选择并行度以匹配线程池大小,但这将确保通过f1 的平均次数将等于通过f2 的平均次数。

【讨论】:

  • 我给出的答案只是以 Akka 流为例——你可以使用任何流技术来做到这一点。您不会认真地期望人们不会仅仅因为他们没有在最初的问题中提到他们想要使用这样的抽象而使用为完美处理他们的用例而构建的抽象吗?
  • 取决于“抽象”是标准(“香草”)scala 还是重型附加库。这不是我“期望”的问题,而是问题实际提出的问题。但是,是的,我不希望人们认真考虑将整个 akka 引入他们的架构中,只是为了对一堆未来进行排序。
【解决方案3】:

这有什么问题?

listOfInts.par.map(f1 andThen f2)

【讨论】:

    【解决方案4】:

    从异步计算的角度来看,从问题中不清楚是否需要将 f1 和 f2 分离为 processA 和 processB - f1 的结果始终且仅传递给 f2,可以这样做作为单个计算“f1 andThen f2”的一部分,并且在单个 Future 中更简单。

    如果是这种情况,那么问题就简化为“如何在可能很大的输入上运行异步计算,从而限制正在运行的、衍生的 Futures”:

    import scala.concurrent.Future
    import java.util.concurrent.Semaphore
    import scala.concurrent.ExecutionContext.Implicits.global
    
    val f: (Int) => Int = i => f2(f1(i))
    
    def process(concurrency: Int, input: List[Int], f: Int => Int): Future[List[Int]] = {
      val semaphore = new Semaphore(concurrency)
      Future.traverse(input) { i =>
        semaphore.acquire()
        Future(f(i)).andThen { case _ => semaphore.release() }
      }
    }
    

    【讨论】:

    • 恐怕这不仅会在不使用给定(在本例中为global)ExecutionContext 上的 BlockContext 的情况下阻塞线程,而且还会有使整个线程池死锁的风险。
    • 它将阻塞一个线程,即执行 process() 的线程。如果要求确实是围绕限制生成的 Futures 数量,则必须在某处进行一些同步“是否已经有最大允许的飞行中异步计算正在进行?”在产生一个新的之前。不过我没有看到僵局——希望能提供一些细节!
    • @ectrop 您可以在获取信号量时阻塞 ExecutionContext 线程(假设 process 也在相同的执行上下文上运行),然后让池中的线程有机会进入 andThen 并释放锁。使用阻塞可能会有所帮助,但我认为仍有可能出现死锁。
    • 信号量获取只会阻塞一个线程——总是同一个线程——运行 process()。如果这是 EC 池中唯一的线程,或者如果所有其他线程都被阻塞,那么可以肯定 - 这是一个死锁。但是在合理的配置下,如果您有足够的(甚至是一个)其他线程可用于运行 f1 和 f2 计算,那么这些线程也可以用于释放信号量。真的不知道死锁是从哪里来的。
    • 相同的执行上下文被隐式传递给Future.traverseFuture.applyandThentraverse 将以 EC 线程为代价运行主体(包括信号量获取)。如果在任何线程有机会释放锁之前将所有 EC 线程分配给信号量获取,会发生什么?
    【解决方案5】:

    可能是这样的:

    items.foreach { it => processB(Seq(it)).onComplete(...) }
    

    当然,如果您的f2f1 重得多,这也无济于事。在这种情况下,您将需要一些更明确的协调:

    val batchSize = 10
    val sem = new Semaphore(batchSize) 
    items.foreach { it => 
       sem.acquire
       processB(Seq(It))
        .andThen { case _ => sem.release }
        .onComplete { ... }
    }
    

    【讨论】:

    • 我觉得信号量的想法很有趣。保持将其与堆栈中的 onFailure 内容正确集成。
    • .andThen 也会处理故障。不应该是“整合”所需的任何其他东西
    • 不需要引入不必要的阻塞。期货完美地描述了时间上的依赖关系。
    • @ViktorKlang 如果你有办法在不阻塞的情况下解决它,请告诉。您的解决方案(1)也有阻塞(仅仅因为您将其隐藏在 flatMap 中并不意味着它不会发生,并且(2)不如这个,尽管更复杂:它等待整个完成在开始下一个之前的前一个块。如果不同任务的执行时间波动很大,这会导致效率低下,并导致突发负载和空闲资源周期。
    • 我认为如果主线程来自同一个 ExecutionContext 隐式传递给andThen 你可能会死锁(与@ectro 答案相同)。顺便说一句,@ViktorKlang 是 scala 期货库的作者之一。 (你好像没有意识到这一点)
    猜你喜欢
    • 1970-01-01
    • 2010-10-15
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-03-24
    • 2011-01-31
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多