中断延迟迭代的功能方式取决于超时和上一个和下一个之间的比较,而 LazyList 与 Stream

Cer*_*vEd 1 functional-programming scala

背景

我有以下场景。我想重复执行来自外部库的类的方法,并且我想这样做直到满足某个超时条件和结果条件(与之前的结果相比)。此外,我想收集返回值,即使在“失败”运行(带有“失败”结果条件的运行,应该中断进一步执行)。

到目前为止,我已经通过初始化一个 empty var result: Result, avar stop: Boolean并使用一个while在条件为真时运行的循环并修改外部状态来实现这一点。我想摆脱这个并使用功能方法。

一些上下文。每次运行预计运行 0 到 60 分钟,迭代的总时间上限为 60 分钟。理论上,这段时间内执行的次数没有限制,但在实践中,一般是2-60次。

问题是,运行需要很长时间,所以我需要停止执行。我的想法是使用某种懒惰IteratorStream加上scanLeftOption

代码

锅炉板

此代码不是特别相关,但用于我的方法示例并提供相同但有些随机的伪运行时结果。

import scala.collection.mutable.ListBuffer
import scala.util.Random
val r = Random
r.setSeed(1)
val sleepingTimes: Seq[Int] = (1 to 601)
  .map(x => Math.pow(2, x).toInt * r.nextInt(100))
  .toList
  .filter(_ > 0)
  .sorted
val randomRes = r.shuffle((0 to 600).map(x => r.nextInt(10)).toList)
case class Result(val a: Int, val slept: Int)
class Lib() {
  def run(i: Int) = {
    println(s"running ${i}")
    Thread.sleep(sleepingTimes(i))
    Result(randomRes(i), sleepingTimes(i))
  }
}
case class Baz(i: Int, result: Result)
val lib = new Lib()
val timeout = 10 * 1000
Run Code Online (Sandbox Code Playgroud)

虽然接近

val iteratorStart = System.currentTimeMillis()
val iterator = for {
  i <- (0 to 600).iterator
  if System.currentTimeMillis() < iteratorStart + timeout
  f = Baz(i, lib.run(i))
} yield f
val iteratorBuffer = ListBuffer[Baz]()
if (iterator.hasNext) iteratorBuffer.append(iterator.next())
var run = true
while (run && iterator.hasNext) {
  val next = iterator.next()
  run = iteratorBuffer.last.result.a < next.result.a
  iteratorBuffer.append(next)
}
Run Code Online (Sandbox Code Playgroud)

流方法 (Scala.2.12)

完整示例

val streamStart = System.currentTimeMillis()
val stream = for {
  i <- (0 to 600).toStream
  if System.currentTimeMillis() < streamStart + timeout
} yield Baz(i, lib.run(i))
var last: Option[Baz] = None
val head = stream.headOption
val tail = if (stream.nonEmpty) stream.tail else stream
val streamVersion = (tail
  .scanLeft((head, true))((x, y) => {
    if (x._1.exists(_.result.a > y.result.a)) (Some(y), false)
    else (Some(y), true)
  })
  .takeWhile {
    case (baz, continue) =>
      if (!baz.eq(head)) last = baz
      continue
  }
  .map(_._1)
  .toList :+ last).flatten
Run Code Online (Sandbox Code Playgroud)

LazyList 方法 (Scala 2.13)

完整示例

val lazyListStart = System.currentTimeMillis()
val lazyList = for {
  i <- (0 to 600).to(LazyList)
  if System.currentTimeMillis() < lazyListStart + timeout
} yield Baz(i, lib.run(i))
var last: Option[Baz] = None
val head = lazyList.headOption
val tail = if (lazyList.nonEmpty) lazyList.tail else lazyList
val lazyListVersion = (tail
  .scanLeft((head, true))((x, y) => {
    if (x._1.exists(_.result.a > y.result.a)) (Some(y), false)
    else (Some(y), true)
  })
  .takeWhile {
    case (baz, continue) =>
      if (!baz.eq(head)) last = baz
      continue
  }
  .map(_._1)
  .toList :+ last).flatten
Run Code Online (Sandbox Code Playgroud)

结果

这两种方法似乎都能产生正确的最终结果:

List(Baz(0,Result(4,170)), Baz(1,Result(5,208)))
Run Code Online (Sandbox Code Playgroud)

并且它们会根据需要中断执行。

编辑:期望的结果是不执行下一次迭代,但仍返回导致中断的迭代结果。因此想要的结果是

List(Baz(0,Result(4,170)), Baz(1,Result(5,208)), Baz(2,Result(2,256))
Run Code Online (Sandbox Code Playgroud)

并且lib.run(i)应该只运行 3 次。

这是通过while方法以及方法实现的,LazyList但不是Stream执行lib.run4 次的方法(糟糕!)。

有没有另一种无状态方法,希望它更优雅?

编辑

我意识到我的示例是错误的并且没有返回它应该返回的“失败”结果,并且它们在停止条件之外继续执行。我重写了代码和示例,但我相信问题的精神是一样的。

Lui*_*rez 5

我会使用更高级别的东西,比如fs2
(或任何其他高级流媒体库,例如:monix observablesakka 流zio zstreams

def runUntilOrTimeout[F[_]: Concurrent: Timer, A](work: F[A], timeout: FiniteDuration)
                                                 (stop: (A, A) => Boolean): Stream[F, A] = {
  val interrupt =
    Stream.sleep_(timeout)

  val run =
    Stream
      .repeatEval(work)
      .zipWithPrevious                                         
      .takeThrough {
        case (Some(p), c) if stop(p, c) => false
        case _                          => true
      } map {
        case (_, c) => c
      }

  run mergeHaltBoth interrupt
}
Run Code Online (Sandbox Code Playgroud)

你可以看到它在这里工作。