【问题标题】:How to control the concurrency of future.sequence in scala?如何控制scala中future.sequence的并发?
【发布时间】:2018-09-30 04:26:47
【问题描述】:

我知道我可以通过将Seq[Future[T]] 转换为Future[Seq[T]]

  val seqFuture = Future.sequence(seqOfFutures)
  seqFuture.map((seqT: Seq[T]) => {...})

我现在的问题是,我在该序列中有 700 个期货,我希望能够控制其中有多少是并行解决的,因为每个期货将调用内部休息 api,并同时有 700 个请求就像对那个服务器发起一个dos攻击。

我宁愿一次只解决 10 个未来。

我怎样才能做到这一点?


尝试pamu's answer我看到错误:

[error] /home/philipp/src/bluebat/src/main/scala/com/dreamlines/metronome/service/JobFetcher.scala:32:44: com.dreamlines.commons.LazyFuture[A] does not take parameters
[error]         val batch = Future.sequence(c.map(_()))
[error]                                            ^
[error] /home/philipp/src/bluebat/src/main/scala/com/dreamlines/metronome/service/JobFetcher.scala:32:28: no type parameters for method sequence: (in: M[scala.concurrent.Future[A]])(implicit cbf: scala.collection.generic.CanBuildFrom[M[scala.concurrent.Future[A]],A,M[A]], implicit executor: scala.concurrent.ExecutionContext)scala.concurrent.Future[M[A]] exist so that it can be applied to arguments (List[Nothing])
[error]  --- because ---
[error] argument expression's type is not compatible with formal parameter type;
[error]  found   : List[Nothing]
[error]  required: ?M[scala.concurrent.Future[?A]]
[error]         val batch = Future.sequence(c.map(_()))
[error]                            ^
[error] /home/philipp/src/bluebat/src/main/scala/com/dreamlines/metronome/service/JobFetcher.scala:32:42: type mismatch;
[error]  found   : List[Nothing]
[error]  required: M[scala.concurrent.Future[A]]
[error]         val batch = Future.sequence(c.map(_()))
[error]                                          ^
[error] /home/philipp/src/bluebat/src/main/scala/com/dreamlines/metronome/service/JobFetcher.scala:32:36: Cannot construct a collection of type M[A] with elements of type A based on a collection of type M[scala.concurrent.Future[A]].
[error]         val batch = Future.sequence(c.map(_()))
[error]                                    ^
[error] four errors found

【问题讨论】:

标签: scala asynchronous concurrency future


【解决方案1】:

向左折叠

简单的foldLeft 可用于控制一次并发运行的期货数量。

首先,让我们创建一个名为LazyFuture的案例类

case class LazyFuture[+A](f: Unit => Future[A]) {
  def apply() = f()
}

object LazyFuture {
  def apply[A](f: => A)(implicit ec: ExecutionContext): LazyFuture[A] = LazyFuture(_ => Future(f))

  def apply[A](f: => Future[A])(implicit ec: ExecutionContext): LazyFuture[A] = LazyFuture(_ => f)
}

LazyFuture 立即停止未来运行

val list: List[LazyFuture[A]] = ...


list.grouped(concurFactor).foldLeft(Future.successful(List.empty[A])){ (r, c) =>
  val batch = Future.sequence(c.map(_()))
  batch.flatMap(values => r.map(rs => rs ++ values))
}

相应地更改concurFactor 以同时运行多个期货。

concurFactor of 1 将同时运行一个未来

concurFactor of 2 将同时运行两个期货

等等……

def executeBatch[A](list: List[LazyFuture[A]])(concurFactor: Int) =
   list.grouped(concurFactor).foldLeft(Future.successful(List.empty[A])){ (r, c) =>
      val batch = Future.sequence(c.map(_()))
      r.flatMap(rs => batch.map(values => rs ++ values))
    }

完整代码

  case class LazyFuture[+A](f: Unit => Future[A]) {
    def apply() = f()
  }

  object LazyFuture {
    def apply[A](f: => A)(implicit ec: ExecutionContext): LazyFuture[A] = LazyFuture(_ => Future(f))

    def apply[A](f: => Future[A])(implicit ec: ExecutionContext): LazyFuture[A] = LazyFuture(_ => f)
  }

  def executeBatch[A](list: List[LazyFuture[A]])(concurFactor: Int)(implicit ec: ExecutionContext): Future[List[A]] =
    list.grouped(concurFactor).foldLeft(Future.successful(List.empty[A])) { (r, c) =>
      val batch = Future.sequence(c.map(_ ()))
      r.flatMap(rs => batch.map(values => rs ++ values))
    }

限制执行上下文

您还可以通过限制执行池中的线程数来限制计算资源。但是,这个解决方案不是那么灵活。就个人而言,我不喜欢它。

val context: ExecutionContext = 
  ExecutionContext.fromExecutor(Executors.newFixedThreadPool(8))

您必须记住传递正确的执行上下文,这是一个隐式值。有时我们不知道哪个隐式在范围内。有问题

警告

当未来构造如下时

val foo = Future {
     1 + 2
} // future starts executing

LazyFuture(foo) // Not a right way

foo已经开始执行,无法控制。

构造LazyFuture的正确方法

val foo = LazyFuture {
  1 + 2
}

val foo = LazyFuture {
  Future {
   1 + 2
  }
}

工作示例

package main

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

object Main {

  case class LazyFuture[A](f: Unit => Future[A]) {
    def apply(): Future[A] = f()
  }

  object LazyFuture {
    def apply[A](f: => A)(implicit ec: ExecutionContext): LazyFuture[A] = LazyFuture(_ => Future(f))
    def apply[A](f: => Future[A]): LazyFuture[A] = LazyFuture(_ => f)
  }

  def executeBatch[A](list: List[LazyFuture[A]])(concurFactor: Int)
    (implicit ec: ExecutionContext): Future[List[A]] =
    list.grouped(concurFactor).foldLeft(Future.successful(List.empty[A])) { (r, c) =>
      val batch = Future.sequence(c.map(_ ()))
      r.flatMap(rs => r.map(values=> rs ++ values))
    }

  def main(args: Array[String]): Unit = {
    import scala.concurrent.ExecutionContext.Implicits.global


    val futures: Seq[LazyFuture[Int]] = List(1, 2, 3, 4, 5).map { value =>
      LazyFuture {
        println(s"value: $value started")
        Thread.sleep(value * 200)
        println(s"value: $value stopped")
        value
      }
    }
    val f = executeBatch(futures.toList)(2)
    Await.result(f, Duration.Inf)
  }

}

【讨论】:

  • 限制执行池线程数不灵活是什么意思?朝那个方向走我会失去什么?
  • @k0pernikus 您必须记住传递正确的执行上下文,这是一个隐式值。有时我们不知道哪个隐式在范围内。它的越野车
  • 如何将List[Future[T]] 转换为List[LazyFuture[T]]
  • @k0pernikus 添加了一个构造函数。请看一看。警告:Future 必须是未评估的。
  • @pamu 您的解决方案在给定时间无法完全运行n 期货。当您开始处理该组时,它会完全利用该池。但是在准备好几个任务后,它不会从下一组中获取任务。并且该组中的最后一个任务将单独运行。
【解决方案2】:

并发是 Scala 的 Futures 由 ExecutionContext 控制。请注意,期货在创建后立即开始在上下文中执行,因此 Future.sequenceExecutionContext 并不重要。从序列创建原始期货时,您必须提供适当的上下文。

默认上下文ExecutionContext.global(通常通过import scala.concurrent.ExecutionContext.Implicits.global 导入)使用与处理器内核一样多的线程,但它也可以为阻塞任务创建许多额外的线程,这些线程包含在scala.concurrent.blocking 中。这通常是所需的行为,但它不适合您的问题。

幸运的是,您可以使用ExecutionContext.fromExecutor 方法来包装Java 线程池。例如:

import java.util.concurrent.Executors
import scala.concurrent.ExecutionContext

val context = ExecutionContext.fromExecutor(Executors.newFixedThreadPool(10))
val seqOfFutures = Seq.fill(700)(Future { callRestApi() }(context))
val sequenceFuture = Future.sequence(seqOfFutures)(ExecutionContext.global)

当然也可以隐式提供上下文:

implicit val context: ExecutionContext = 
  ExecutionContext.fromExecutor(Executors.newFixedThreadPool(10))
val seqOfFutures = Seq.fill(700)(Future { callRestApi() })
// This `sequence` uses the same thread pool as the original futures
val sequenceFuture = Future.sequence(seqOfFutures) 

【讨论】:

    猜你喜欢
    • 2018-08-13
    • 2020-06-15
    • 2011-01-19
    • 1970-01-01
    • 2013-08-20
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多