【问题标题】:Scala: joining / waiting for growing queue of futuresScala:加入/等待增长的期货队列
【发布时间】:2017-07-30 11:25:31
【问题描述】:

我启动了几个异步进程,如果需要,它们反过来可以启动更多进程(想想遍历目录结构或类似的东西)。每个进程都会返回一些东西,最后我想等待所有这些都完成并安排一个函数来处理结果集合。

天真的尝试

我的解决方案尝试使用了可变的ListBuffer(我不断添加我生成的期货)和Future.sequence 来安排一些函数在完成此缓冲区中列出的所有这些期货时运行。

我准备了一个最小的例子来说明这个问题:

object FuturesTest extends App {
  var queue = ListBuffer[Future[Int]]()

  val f1 = Future {
    Thread.sleep(1000)
    val f3 = Future {
      Thread.sleep(2000)
      Console.println(s"f3: 1+2=3 sec; queue = $queue")
      3
    }
    queue += f3
    Console.println(s"f1: 1 sec; queue = $queue")
    1
  }
  val f2 = Future {
    Thread.sleep(2000)
    Console.println(s"f2: 2 sec; queue = $queue")
    2
  }

  queue += f1
  queue += f2
  Console.println(s"starting; queue = $queue")

  Future.sequence(queue).foreach(
    (all) => Console.println(s"Future.sequence finished with $all")
  )

  Thread.sleep(5000) // simulates app being alive later
}

它首先调度f1f2 期货,然后f3 将在1 秒后以f1 分辨率调度。 f3 本身将在 2 秒内解决。因此,我期望得到的是以下内容:

starting; queue = ListBuffer(Future(<not completed>), Future(<not completed>))
f1: 1 sec; queue = ListBuffer(Future(<not completed>), Future(<not completed>), Future(<not completed>))
f2: 2 sec; queue = ListBuffer(Future(Success(1)), Future(<not completed>), Future(<not completed>))
f3: 1+2=3 sec; queue = ListBuffer(Future(Success(1)), Future(Success(2)), Future(<not completed>))
Future.sequence finished with ListBuffer(1, 2, 3)

但是,我实际上得到了:

starting; queue = ListBuffer(Future(<not completed>), Future(<not completed>))
f1: 1 sec; queue = ListBuffer(Future(<not completed>), Future(<not completed>), Future(<not completed>))
f2: 2 sec; queue = ListBuffer(Future(Success(1)), Future(<not completed>), Future(<not completed>))
Future.sequence finished with ListBuffer(1, 2)
f3: 1+2=3 sec; queue = ListBuffer(Future(Success(1)), Future(Success(2)), Future(<not completed>))

...这很可能是因为我们等待的期货列表在Future.sequence 的初始调用期间是固定的,以后不会更改。

工作,但丑陋的尝试

最终,我用这段代码让它按照我的意愿行事:

  waitForSequence(queue, (all: ListBuffer[Int]) => Console.println(s"finished with $all"))

  def waitForSequence[T](queue: ListBuffer[Future[T]], act: (ListBuffer[T] => Unit)): Unit = {
    val seq = Future.sequence(queue)
    seq.onComplete {
      case Success(res) =>
        if (res.size < queue.size) {
          Console.println("... still waiting for tasks")
          waitForSequence(queue, act)
        } else {
          act(res)
        }
      case Failure(exc) =>
        throw exc
    }
  }

这按预期工作,最终获得所有 3 个期货:

starting; queue = ListBuffer(Future(<not completed>), Future(<not completed>))
f1: 1 sec; queue = ListBuffer(Future(<not completed>), Future(<not completed>), Future(<not completed>))
f2: 2 sec; queue = ListBuffer(Future(Success(1)), Future(<not completed>), Future(<not completed>))
... still waiting for tasks
f3: 1+2=3 sec; queue = ListBuffer(Future(Success(1)), Future(Success(2)), Future(<not completed>))
finished with ListBuffer(1, 2, 3)

但还是很丑。它只是重新启动Future.sequence 等待,如果它看到在完成时队列长于结果数,希望下次完成时情况会更好。当然,这很糟糕,因为它会耗尽堆栈,并且如果此检查将在创建未来和将其附加到队列之间的一个小窗口中触发,则可能容易出错。


是否可以在不使用 Akka 重写所有内容或诉诸 Await.result 的情况下这样做(I can't actually use,因为我的代码是为 Scala.js 编译的)。

【问题讨论】:

    标签: multithreading scala future scala.js


    【解决方案1】:

    就像 Justin 提到的,您不能丢失对在其他期货中生成的期货的引用,您应该使用 map 和 flatMap 将它们链接起来。

    val f1 = Future {
      Thread.sleep(1000)
      val f3 = Future {
        Thread.sleep(2000)
        Console.println(s"f3: 1+2=3 sec")
        3
      }
      f3.map{
        r =>
          Console.println(s"f1: 1 sec;")
          Seq(1, r)
      }
    }.flatMap(identity)
    
    val f2 = Future {
      Thread.sleep(2000)
      Console.println(s"f2: 2 sec;")
      Seq(2)
    }
    
    val futures = Seq(f1, f2)
    
    Future.sequence(futures).foreach(
      (all) => Console.println(s"Future.sequence finished with ${all.flatten}")
    )
    
    Thread.sleep(5000) // simulates app being alive later
    

    这适用于最小的示例,我不确定它是否适用于您的实际用例。结果是:

    f2: 2 sec;
    f3: 1+2=3 sec
    f1: 1 sec;
    Future.sequence finished with List(1, 3, 2)
    

    【讨论】:

    【解决方案2】:

    执行此操作的正确方法可能是编写您的 Futures。具体来说,f1 不应该只是启动 f3,它可能应该在它上面进行 flatMap——也就是说,f1 的 Future 直到 f3 解决才解决。

    请记住,Future.sequence 是一种备用选项,仅在 Future 都真正断开连接时使用。在您所描述的情况下,存在真正的依赖关系,这些在您实际返回的 Futures 中得到最好的体现。使用 Futures 时,flatMap 是您的朋友,应该是您最先使用的工具之一。 (通常但不总是像for 理解。)

    可以肯定地说,如果您想要一个可变的 Futures 队列,代码结构不正确,但有更好的方法可以做到这一点。特别是在 Scala.js 中(这是我的大部分代码所在的地方,并且非常重未来),我使用这些 Future 的理解不断 - 我认为这是唯一理智的操作方式...

    【讨论】:

      【解决方案3】:

      我不会涉及Future.sequence:它使操作并行化,您似乎正在寻找顺序异步执行。此外,您可能不需要在定义后立即开始期货。组合应如下所示:

      def run[T](queue: List[() => Future[T]]): Future[List[T]] = {
        (Future.successful(List.empty[T]) /: queue)(case (f1, f2) =>
        f1() flatMap (h => )
        )
      
      val t0 = now
      
      def f(n: Int): () => Future[String] = () => {
        println(s"starting $n")
        Future[String] {
          Thread.sleep(100*n)
          s"<<$n/${now - t0}>>"
        }
      }
      
      println(Await.result(run(f(7)::f(10)::f(20)::f(3)::Nil), 20 seconds))
      

      诀窍是不要过早推出期货;这就是为什么我们有 f(n) 直到我们用 () 调用它才会开始。

      【讨论】:

        猜你喜欢
        • 2015-06-03
        • 2020-10-31
        • 2020-06-12
        • 2013-08-02
        • 1970-01-01
        • 1970-01-01
        • 2015-08-24
        • 2018-08-22
        • 2023-03-16
        相关资源
        最近更新 更多