我有一个分页的资源,我想与Monix递归使用它。我想要一个Observable,它将发出下载的元素并递归使用页面。这是一个简单的例子。当然不行。它发出第一页,然后是第一页+第二页,然后是第一+第二+第三页。我希望它先发出,然后发出第二,然后发出第三,依此类推。
object Main extends App {
sealed trait Event
case class Loaded(xs: Seq[String]) extends Event
// probably should just finish stream instead of this event
case object Done extends Event
// here is the problem
def consume(page: Int, size: Int):Observable[Event] = {
Observable.fromFuture(getPaginatedResource(page, size)).concatMap{ xs =>
if (xs.isEmpty) Observable.pure(Done)
else Observable.concat(Observable.pure(Loaded(xs)), consume(page + 1, size + 5))
}
}
def getPaginatedResource(page: Int, size: Int):Future[Seq[String]] = Future {
if (page * size > 100) Seq.empty
else 0 to size map (x => s"element $x")
}
consume(page = 0, size = 5).foreach(println)
}
Run Code Online (Sandbox Code Playgroud)
有任何想法吗?
UPD
对不起,似乎正在运行,并且我只有一个bug size + 5。看来问题已经解决,但是如果您发现我做错了,请告诉我。
通常建议在使用时尽可能避免递归Observable。因为它不容易可视化,而且通常更容易出错。
一个想法是使用,scanEvalF因为它会在每一步上发出项目。
sealed trait Event
object Event {
case class Loaded(page: Int, size: Int, items: Seq[String]) extends Event
}
def getPaginatedResource(page: Int, size: Int): Task[Loaded] = Task.pure {
if (page * size > 100) Loaded(page, size, Seq.empty)
else Loaded(page, size, 0.to(size).map(x => s"element $x"))
}
def consume(page: Int, size: Int): Observable[Event] = {
Observable
.interval(0.seconds)
.scanEvalF(getPaginatedResource(page, size)) { (xs, _) =>
getPaginatedResource(xs.page + 1, xs.size + 5)
} // will emit items on each step
.takeWhileInclusive(_.items.nonEmpty) // only take until list is empty
}
consume(0, 5)
.foreachL(println)
.runToFuture
Run Code Online (Sandbox Code Playgroud)
| 归档时间: |
|
| 查看次数: |
150 次 |
| 最近记录: |