【问题标题】:Scalaz stream group sorted database resultsScalaz 流组排序的数据库结果
【发布时间】:2013-09-30 20:38:34
【问题描述】:

我在我的代码中看到了一个常见的模式。我已经对数据库中的结果进行了排序,我需要将它们以嵌套结构发出。我希望它可以流式传输,因此我希望一次在内存中拥有尽可能少的记录。使用 TravesableLike.groupBy 假定数据未排序,因此它不必要地填充可变映射。我想保持这种真正的流媒体。 scalaz-stream 在这里有用吗?

val sql = """select grandparent_id, parent_id, child_id
  from children
  where grandparent_id = ?
  order by grandparent_id, parent_id, child_id"""

def elementsR[P, R](invoker: Invoker[P, R], param: P): Process[Task, R] =
  // Invoker.elements returns trait CloseableIterator[+T] extends Iterator[T] with Closeable
  resource(Task.delay(invoker.elements(param)))(
    src => Task.delay(src.close)) { src =>
      Task.delay { if (src.hasNext) src.next else throw End }
  }

def dbWookie {
  // grandparent_id, (grandparent_id, parent_id, child_id)
  val invoker = Q.query[Int, (Int, Int, Int)](sql)
  val es = elementsR(invoker, 42)

  // ?, ?, ?

  // nested emits (42, ((35, (1, 3, 7)), (36, (8, 9, 12))))
}

我没有在 Process 上看到太多函数,如 foldLeft 和 scanLeft,所以我不确定如何检测 grandparent_id、parent_id 或 child_id 何时更改并发出组。有什么想法吗?

【问题讨论】:

    标签: scala scalaz scalaz-stream


    【解决方案1】:

    我认为您想要一些与chunkBy 类似的东西。每当谓词函数的结果从true 翻转到false 时,chunkBy 就会发出一个块。

    您可以将其从比较布尔值概括为比较输入的某个任意函数的结果。因此,只要应用于输入的此函数的值发生变化,您就会有一个进程发出一个块:

    def chunkOn[I, A](f: I => A): Process1[I, Vector[I]] = {
      def go(acc: Vector[I], last: A): Process1[I,Vector[I]] =
        await1[I].flatMap { i =>
          val cur = f(i)
          if (cur != last) emit(acc) then go(Vector(i), cur)
          else go(acc :+ i, cur)
        } orElse emit(acc)
      await1[I].flatMap(i => go(Vector(i), f(i)))
    }
    

    REPL 中的快速脏测试,使用 Identity monad 立即强制评估:

    scala> import scalaz.stream._, scalaz.Id._
    import scalaz.stream._
    import scalaz.Id._
    
    scala> val rows = Seq(('a, 'b, 'c), ('a, 'b, 'd), ('b, 'a, 'c), ('b, 'd, 'a))
    rows: Seq[(Symbol, Symbol, Symbol)] = List(('a,'b,'c), ('a,'b,'d), ('b,'a,'c), ('b,'d,'a))
    
    scala> val process = Process.emitSeq[Id, (Symbol, Symbol, Symbol)](rows)
    process: scalaz.stream.Process[scalaz.Id.Id,(Symbol, Symbol, Symbol)] =
      Emit(List(('a,'b,'c), ('a,'b,'d), ('b,'a,'c), ('b,'d,'a)),Halt(scalaz.stream.Process$End$))
    
    scala> process |> chunkOn(_._1)
    res4: scalaz.stream.Process[scalaz.Id.Id,scala.collection.immutable.Vector[(Symbol, Symbol, Symbol)]] =
      Emit(List(Vector(('a,'b,'c), ('a,'b,'d))),Emit(List(Vector(('b,'a,'c), ('b,'d,'a))),Halt(scalaz.stream.Process$End$)))
    

    正如您所建议的,chunkWhen 对当前值和最后一个值使用谓词,并在计算结果为 false 时发出一个块。

    def chunkWhen[I](f: (I, I) => Boolean): Process1[I, Vector[I]] = {
      def go(acc: Vector[I]): Process1[I,Vector[I]] =
        await1[I].flatMap { i =>
          acc.lastOption match {
            case Some(last) if ! f(last, i) => emit(acc) then go(Vector(i))
            case _ => go(acc :+ i)
          }
        } orElse emit(acc)
      go(Vector())
    }
    

    试一试:

    scala> process |> chunkWhen(_._1 == _._1)
    res0: scalaz.stream.Process[scalaz.Id.Id,Vector[(Symbol, Symbol, Symbol)]] =
      Emit(List(Vector(('a,'b,'c), ('a,'b,'d))),Emit(List(Vector(('b,'a,'c), ('b,'d,'a))),Halt(scalaz.stream.Process$End$)))
    

    【讨论】:

    • 好的,很酷。为什么chunkBy的最后一行发生了变化? go(Vector(), false) 变成 await1[I].flatMap(i => go(Vector(i), f(i)))?
    • 我们不能直接调用go——我们没有A类型的值作为last传入,我们也不能发明一个,因为我们没有知道A 是什么类型。而在chunkBy 中,false 可以作为起始值。
    • 你对被称为 chunkWhen() 的东西有什么看法?我们能否将 A 移除并仅依赖 chunkWhen[I](f: (I, I) => Boolean)。还是不够?
    • 我也应该问,在 chunckWhen 的基本情况下该怎么办? await1[I].flatMap(i => go(Vector(i), f(i))) 变成了什么?
    • 我在答案中添加了chunkWhen 的可能实现。我在累加器上使用lastOption 来简化基本情况,尽管您也可以显式传递last 并执行类似于chunkOn 的操作。
    猜你喜欢
    • 2020-05-06
    • 2023-03-19
    • 1970-01-01
    • 1970-01-01
    • 2023-03-13
    • 2013-12-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多