【问题标题】:Building a LazyList with Futures in a non-blocking and lazy way以非阻塞和惰性的方式构建带有 Futures 的 LazyList
【发布时间】:2020-10-06 02:55:27
【问题描述】:

我正在为传递数据块的 Clojure 通道构建 Scala 外观,我想将其表示为 LazyList[Future[Either[String, Int]]],其中左侧可以保存错误消息,右侧可以保存数据。从 Channel 检索每个块是一个阻塞操作,因此我想将每个块封装在 Future 中。

每个块结果类型决定了我们应该如何继续构建惰性列表:

  1. null:频道上没有更多结果,返回列表
  2. 字符串:添加Left(error)并返回列表
  3. Int:添加 Right(data) 并递归下一个块

我的问题是我们是否可以以一种惰性且非阻塞的方式构建这样的列表?

这是我到目前为止想出的,但是头部被评估(不是懒惰的)并且 Await.result 块:

// Clojure "Channel" dummy
case class Channel(vs: Any*) {
  private val it = vs.toIterable.iterator

  // equivalent to the `<!!` Clojure function
  def chunk: Future[Any] = Future {
    // This imitates an expensive blocking operation
    if (it.hasNext) {
      val value = it.next
      println("Retrieving value: " + value)
      value
    } else {
      null
    }
  }
}

def lazyList(channel: Channel): LazyList[Future[Either[String, Int]]] = {
  val ll = channel.chunk.map {
    case null          => LazyList.empty[Future[Either[String, Int]]] // No more values
    case error: String => Future(Left(error)) #:: LazyList.empty[Future[Either[String, Int]]]
    case data: Int     => Future(Right(data)) #:: lazyList(channel)
  }
  Await.result(ll, Duration.Inf)
}

val ll = lazyList(Channel(0, 1, "error"))
// Retrieving value: 0
ll(0)
// (no output since value 0 has already been calculated and memoized)
ll(1)
// Retrieving value: 1
ll(2)
// Retrieving value: error

我希望看到的是:

val ll2 = lazyList2(Channel(0, 1, "error"))
// (no computation)
ll2(0)
// Retrieving value: 0
ll2(1)
// Retrieving value: 1
ll2(2)
// Retrieving value: error

【问题讨论】:

  • 您是否认为 fs2 是一种流媒体解决方案,而不是您自己的解决方案?它的设计非常好,可以有效且正确地解决一些类似 Future 的事物产生数据流的问题。
  • 你不能有一个 LazyList[A] 具有由 Future 产生的值并且不会阻塞。当你想产生下一个 A 并且有一个 Future 产生 A 时,你必须阻塞等待它。您的阻塞实现将按照您希望在 scala 2.13.4 中看到的那样工作,它修复了对 #:: 的热切评估

标签: scala nonblocking


【解决方案1】:

如果您使用的是 fs2,则可以从通道构建流。给定一个函数

def nextChunk: Future[A] = ???

你可以构建一个流

val myStream: Stream[IO, A] = Stream.eval(IO.fromFuture(IO(nextChunk))).repeat

在您的具体示例中,您的 AAny,您知道在运行时是 IntStringnull。您可以先将其提升为 Option[Either[String, Int]]

def typedChunk(channel: Channel): IO[Option[Either[String, Int]]] = 
  IO.fromFuture(IO(channel.nextChunk)).map {
    case null      => None
    case s: String => Some(Left(s))
    case i: Int    => Some(Right(i))
  }

然后您可以构建流,以None 终止

def myTerminatedStream(channel: Channel): Stream[IO, Either[String, Int]] = 
  Stream.eval(typedChunk(channel)).repeat.unNoneTerminate

这完成了保持引用透明度并确保它为您提供正确评估语义的所有艰苦工作。

您使用 LazyList 请求的语义将是棘手的:您只会在 Future 完成评估后知道您的块为空,因此您需要评估 Future 以了解您的列表是否为空。 LazyList 能够做到这一点,但只能来自阻塞操作,而不是来自 Future。

【讨论】:

  • 谢谢,@Martijn。一切都是懒惰的,太棒了。需要更多地探索 fs2 以了解我将如何将它用于我的外观。当我有一个示例启动并运行以供将来参考时,我会回来。
猜你喜欢
  • 2020-12-07
  • 2013-10-29
  • 1970-01-01
  • 2011-07-14
  • 1970-01-01
  • 1970-01-01
  • 2016-08-26
  • 2018-05-29
  • 1970-01-01
相关资源
最近更新 更多