【问题标题】:scala Future to run sequential jobsscala Future 运行顺序作业
【发布时间】:2018-02-21 11:42:24
【问题描述】:

我正在尝试按顺序启动三个作业,但是当我尝试此代码时:

val jobs = Seq("stream.Job1","stream.Job2","stream.Job3")
    Future.sequence {
          jobs.map { jobClass =>
            Future {
              println(s"Starting the spark job from class $jobClass...")
              % gcloud("sparkC", "jobs", "submit", "spark", s"--cluster=$clusterName", s"--class=$jobClass", "--region=global", s"--jars=$JarFile")
              println(s"Starting the spark job from class $jobClass...DONE")
            }
          }
        }  

我同时获得这三个工作,然后是顺序的。 我认为解决方案是使用flatMap 但我无法实现它。
请帮忙。

【问题讨论】:

  • “并行,然后顺序”是什么意思?作业是否在完成一半时就停止了,然后将自己排成某种顺序队列?
  • 您是否尝试运行依赖于输出三个作业的作业(并行运行)?
  • 如果你想要顺序执行,你为什么还要使用 Future?
  • 我想逐个工作(一个接一个)。当job1结束时job2开始

标签: scala future


【解决方案1】:

试试这个

val jobs = Seq("stream.Job1","stream.Job2","stream.Job3")
jobs.foldLeft(Future.successful[Unit]()) {
  case (result, jobClass) =>
    result.flatMap[Unit] {_ =>
      Future {
        println(s"Starting the spark job from class $jobClass...")
        % gcloud("sparkC", "jobs", "submit", "spark", s"--cluster=$clusterName", s"--class=$jobClass", "--region=global", s"--jars=$JarFile")
        println(s"Starting the spark job from class $jobClass...DONE")
      }
    }.
      recoverWith {
      case NonFatal(e) => result
    }
}

这将遍历您的作业,并在上一个完成后立即运行下一个未来。如果其中任何一个失败,我添加了recoverWith 块以独立处理所有Futures

【讨论】:

  • 感谢您的回复,您可以使用我的示例编辑您的代码吗?
【解决方案2】:

如果作业不相互依赖,并且如果您想要一个结果列表 最后,你可以使用这个:

import scala.concurrent._
def runIndependentSequentially[X]
  (futs: List[() => Future[X]])
  (implicit ec: ExecutionContext): Future[List[X]] = futs match {
  case Nil => Future { Nil }
  case h :: t => for {
    x <- h()
    xs <- runIndependentSequentially(t)
  } yield x :: xs
}

现在您可以在未来的工作列表中使用它,如下所示:

import scala.concurrent.ExecutionContext.Implicits.global
import scala.concurrent.duration._
import scala.language.postfixOps

val jobs = List("stream.Job1","stream.Job2","stream.Job3")
val futFactories = jobs.map { jobClass =>
  () => Future {
    println(s"Starting the spark job from class $jobClass...")
    Thread.sleep(5000)
    "result[" + jobClass + "," + (System.currentTimeMillis / 1000) % 3600 + "]"
  }
}

println(Await.result(runIndependentSequentially(futFactories), 30 seconds))

这会产生以下输出:

Starting the spark job from class stream.Job1...
Starting the spark job from class stream.Job2...
Starting the spark job from class stream.Job3...
List(result[stream.Job1,3011], result[stream.Job2,3016], result[stream.Job3,3021])

UPDATE:用List[() =&gt; Future[X]] 替换期货列表,这样 甚至在参数传递给 runIndependentSequentially 方法。非常感谢@Evgeny 指出!

【讨论】:

  • 恐怕这行不通。例如,在输出中添加时间戳,您将看到所有期货将在映射后立即启动,而不是按顺序启动。
  • @Evgeny 谢谢,已修复!
猜你喜欢
  • 2021-05-23
  • 2019-03-30
  • 1970-01-01
  • 2016-03-26
  • 2018-08-01
  • 2020-04-29
  • 1970-01-01
  • 2019-03-27
  • 1970-01-01
相关资源
最近更新 更多