【问题标题】:Lazy Pagination in Scala (Stream/Iterator of Iterators?)Scala中的延迟分页(迭代器的流/迭代器?)
【发布时间】:2021-04-21 20:23:39
【问题描述】:

我通过调用def readPage(pageNumber: Int): Iterator[Record],一次一页从数据库 API 中顺序读取大量记录(每页记录数未知)

我正在尝试将这个 API 封装在 Stream[Iterator[Record]] 或 Iterator[Iterator[Record]] 之类的东西中,以一种功能性的方式,理想情况下没有可变状态,具有恒定的内存占用,以便我可以将其视为无限的页面流或迭代器序列,并从客户端抽象出分页。客户端可以对结果进行迭代,通过调用 next() 它将检索下一页 (Iterator[Record])。

在 Scala 中实现这一点最惯用和最有效的方法是什么。

编辑:需要一次一页地获取和处理记录,无法维护内存中所有页面的所有记录。如果一页失败,则抛出异常。大量页面/记录对于所有实际目的来说意味着无限。我想将其视为无限的页面流(或迭代器),每个页面都是有限数量记录的迭代器(例如,小于

我在 Monix 中查看了 BatchCursor,但它的用途不同。

编辑 2:这是当前版本,使用以下 Tomer 的 answer 作为起点,但使用 Stream 而不是 Iterator。 这允许根据 https://stackoverflow.com/a/10525539/165130 消除尾递归的需要,并且有 O(1) 时间用于流前置 #:: 操作(而如果我们通过 ++ 操作连接迭代器,它将是 O(n))

注意:虽然流被延迟评估,Stream memoization 仍可能导致内存爆炸,并且内存管理得到tricky。从val 更改为def 以在下面的def pages = readAllPages 中定义Stream 似乎没有任何效果

def readAllPages(pageNumber: Int = 0): Stream[Iterator[Record]] = {
   val iter: Iterator[Record] = readPage(pageNumber)
   if (iter.isEmpty)
     Stream.empty
   else
    iter #:: readAllPages(pageNumber + 1)
} 
      
//usage
val pages = readAllPages
for{
    page<-pages
    record<-page
    if(isValid(record))
}
process(record)
 

编辑 3: Tomer 的第二个建议似乎是最好的,它的运行时间和内存占用与上述解决方案类似,但更简洁且容易出错。

val pages = Stream.from(1).map(readPage).takeWhile(_.nonEmpty)

注意:Stream.from(1) 创建一个从 1 开始并递增 1 的流,它在 API docs 中

【问题讨论】:

  • 答案是,取决于... 你想达到什么目标?您是否希望同时阅读所有页面?逐个?如果一页失败会发生什么?你重新开始吗?只重试这个页面?什么是大量记录?所有记录可以一起存在于您机器的内存中吗?如果没有,你必须确保完成一个,直到你得到下一个。
  • 我想一次获取并处理一页记录,无法维护内存中的所有页面。如果一页失败,则抛出异常。由于大量记录对所有实际目的都意味着无限,我想将其视为无限流。
  • 只有当您的客户端可以无限期地连接到服务器时,流媒体才是一种解决方案。然后,您可以通过一个响应流式传输所有内容,或者使用例如websockets 在客户端需要时响应下一批结果。如果您无法与服务器建立一个连接并且使用了分页,那么您必须在服务器中存储状态,即带有游标的数据库连接。这意味着您不能处理太多的客户。这就是为什么通常分页等于带有 LIMIT 和 OFFSET 的单独数据库请求(结果可以在查询之间改变)。
  • 所以您想使用什么主要取决于您的 API 使用者必须如何使用它。
  • 对于创建 Stream,您可以在 Stream、LazyList 或 Iterator 中使用 unfold 方法.但是,正如其他人所指出的,在 REST API 上返回它取决于您的框架。此外,公开迭代器通常不是一个好主意,因为如果您错误地使用它,您将产生错误。最后,unfold 解决方案不会处理错误,如果您真的想要功能更强大的 API,最好为每个页面返回 Either 或 Option 或 List。

标签: scala scala-streams


【解决方案1】:

你可以尝试实现这样的逻辑:

def readPage(pageNumber: Int): Iterator[Record] = ???

@tailrec
def readAllPages(pageNumber: Int): Iterator[Iterator[Record]] = {
  val iter = readPage(pageNumber)
  if (iter.nonEmpty) {
    // Compute on records
    // When finishing computing:
    Iterator(iter) ++ readAllPages(pageNumber + 1)
  } else {
    Iterator.empty
  }
}

readAllPages(0)

一个较短的选项是:

Stream.from(1).map(readPage).takeWhile(_.nonEmpty)

【讨论】:

  • 为了简单起见,我喜欢它,唯一的一点是,它需要返回 Iterator[Iterator[Record]] 而不是 Unit 以便客户端可以迭代结果,通过调用 next 它将检索下一页
  • 谢谢,顺便说一句,您的第一个解决方案不是尾递归(“递归调用不在尾位置”),可以通过添加类似缓冲区的参数使其成为尾递归,但最好使用 Stream因为它完全消除了尾递归的需要,正如我刚刚从这个答案中了解到的那样stackoverflow.com/a/10525539/165130
  • 是的,它有帮助,我会修复它并发布
  • 谢谢!接受了第二个建议,因为它似乎是最好的,在上面编辑了我的问题
猜你喜欢
  • 2017-04-18
  • 1970-01-01
  • 1970-01-01
  • 2011-06-17
  • 1970-01-01
  • 2018-02-16
  • 1970-01-01
  • 2018-10-26
  • 2011-07-16
相关资源
最近更新 更多