【发布时间】:2020-10-06 02:55:27
【问题描述】:
我正在为传递数据块的 Clojure 通道构建 Scala 外观,我想将其表示为 LazyList[Future[Either[String, Int]]],其中左侧可以保存错误消息,右侧可以保存数据。从 Channel 检索每个块是一个阻塞操作,因此我想将每个块封装在 Future 中。
每个块结果类型决定了我们应该如何继续构建惰性列表:
- null:频道上没有更多结果,返回列表
- 字符串:添加
Left(error)并返回列表 - 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