【问题标题】:While Loop with Futures in Scala在 Scala 中使用期货的 While 循环
【发布时间】:2019-07-05 04:18:47
【问题描述】:

我有两种方法:

def getNextJob: Future[Option[Job]]

def process(job: Job): Future[Unit]

我想处理所有作业,直到没有剩余作业为止。

我可以使用Await 来做到这一点,例如

private def process()(implicit ctx: ExecutionContext): Future[Unit] = {
    var job: Option[Job] = Await.result(service.getNextJob, FiniteDuration(2, TimeUnit.SECONDS))
    while(job.isDefined) {
      Await.result(process(job.get), FiniteDuration(2, TimeUnit.SECONDS))
      job = Await.result(service.getNextJob, FiniteDuration(2, TimeUnit.SECONDS))
    }
    Future.successful()
  }

但这很丑陋,并且不能正确使用 Futures。有没有办法以某种方式链接期货来代替它?

【问题讨论】:

  • 你想像你的示例代码那样按照严格的顺序process作业,还是可以按任何顺序执行?
  • @Tim 任何订单都可以,我只需要确保一次运行 1 个
  • 我想我对为什么 process 返回 Future 感到困惑,如果您总是等待它完成,然后再继续下一份工作。如果您不等待它完成,那么可以同时处理多个作业。
  • @Tim 我不打算进入实现,但getNextJob 只是从数据库中拉出一份状态为unprocessed 的工作。 process 处理它并将状态更新为done。如果我并行处理多个,则无法保证(使用我当前的实现)该作业只会运行一次,因为对 getNextJob 的两次调用可能会两次返回相同的作业`
  • 在这种情况下,您可能应该避免从process 返回Future,而是让它进行处理,然后返回Unit。这已经在getNextJob 的Future 中运行,因此无需在第一个Future 中嵌套另一个Future。

标签: scala future


【解决方案1】:
def go()(implicit ctx: ExecutionContext): Future[Unit] =
  getNextJob.flatMap { maybeJob ⇒
    if(maybeJob.isDefined) process(maybeJob.get).flatMap(_ ⇒ go())
    else Future.unit
  }

注意:不是尾递归。

【讨论】:

  • 它实际上根本不是递归的 :) 请注意,Option.get 是“代码气味”,应该避免。你应该改用case Some(job) => process(job).flatMap(_ => go()) ; case None => Future.unit 之类的东西
  • 我接受了@Dima 的建议来回答这个问题,非常感谢!
  • 请注意,由于它是一个未来,所以不是尾递归的函数应该不是问题stackoverflow.com/a/16986416/1507124
【解决方案2】:
def processAll()(implicit ec: ExecutionContext): Future[Unit] =
  getNextJob.flatMap {
    case Some(job) => process(job).flatMap(_ => processAll())
    case None => Future.unit
  }

可能同时处理它们:

def processAll()(implicit ec: ExecutionContext): Future[Unit] =
  getNextJob.flatMap {
    case Some(job) => process(job).zipWith(processAll())((_,_) => ())
    case None => Future.unit
  }

【讨论】:

  • 这将一次处理一个作业,对吗?如果是这样,如何修改它以同时处理作业?
  • @CervEd 补充说。不过警告,我是从内存中写出来的,所以它可能无法编译
  • 谢谢!我使用了Future.sequence { (1 to X).map(processAll(...)) } 之类的东西,其中 X 是我希望一次处理多少个请求,但感觉不正确
猜你喜欢
  • 1970-01-01
  • 2017-09-13
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2013-02-17
  • 2021-07-18
  • 1970-01-01
  • 2018-09-08
相关资源
最近更新 更多