Scalaz中的异步迭代处理

Aar*_*rup 6 asynchronous scala scalaz iterate

我一直在使用Scalaz 7迭代器来处理恒定堆空间中的大型(即无界)数据流.

在代码中,它看起来像这样:

type ErrorOrT[M[+_], A] = EitherT[M, Throwable, A]
type ErrorOr[A] = ErrorOrT[IO, A]

def processChunk(c: Chunk): Result

def process(data: EnumeratorT[Chunk, ErrorOr]): IterateeT[Chunk, ErrorOr, List[Result]] =
  Iteratee.fold[Chunk, ErrorOr, List[Result]](Nil) { (rs, c) =>
    processChunk(c) :: rs
  } &= data
Run Code Online (Sandbox Code Playgroud)

现在我想并行执行处理,一次处理P块数据.我仍然需要限制堆空间,但是假设有足够的堆来存储P块数据和累积的计算结果是合理的.

我知道Task类和想到映射在枚举器上以创建任务流:

data map (c => Task.delay(processChunk(c)))
Run Code Online (Sandbox Code Playgroud)

但我仍然不确定如何管理非决定论.在使用流时,如何尽可能确保P任务正在运行?

第一次尝试:

我首先尝试解决方案是折叠流并创建一个Scala Future来处理每个块.但是,该程序因GC开销错误而爆炸(可能是因为它试图创建所有Futures 时将所有块都拉入内存).相反,当已经存在P任务时,iteratee需要停止消耗输入,并且当任何任务完成时再次恢复.

第二次尝试:

我的下一个尝试是将流分组为P大小的部分,并行处理每个部分,然后在继续下一部分之前加入:

def process(data: EnumeratorT[Chunk, ErrorOr]): IterateeT[Chunk, ErrorOr, Vector[Result]] =
  Iteratee.foldM[Vector[Chunk], ErrorOr, Vector[Result]](Nil) { (rs, cs) =>
    tryIO(IO(rs ++ Await.result(
      Future.traverse(cs) { 
        c => Future(processChunk(c)) 
      }, 
      Duration.Inf)))
  } &= (data mapE Iteratee.group(P))
Run Code Online (Sandbox Code Playgroud)

虽然这不会充分利用可用的处理器(特别是因为处理每个处理器所需的时间Chunk可能差异很大),但这将是一种改进.然而,groupenumeratee似乎泄漏内存-堆的使用情况突然从房顶去.