【问题标题】:scalaz stream structure for growing lists用于增长列表的 scalaz 流结构
【发布时间】:2015-12-07 00:14:51
【问题描述】:

我有一种预感,我可以(应该?)使用 scalaz-streams 来解决我的问题,就像这样。

我有一个起始项 A。我有一个接受 A 并返回 A 列表的函数。

def doSomething(a : A) : List[A]

我有一个以 1 项(起始项)开头的工作队列。当我们处理 (doSomething) 每个项目时,它可能会将许多项目添加到同一个工作队列的末尾。然而,在某些时候(在数百万个项目之后)我们doSomething 上的每个后续项目将开始向工作队列添加越来越少的项目,最终不会添加新项目(doSomething 将为这些项目返回 Nil)。这就是我们知道计算最终会终止的方式。

假设 scalaz-streams 适用于此,请给我一些提示,告诉我应该考虑哪些整体结构或类型来实现它?

一旦完成了使用单个“worker”的简单实现,我还想使用多个 worker 并行处理队列项,例如有一个由 5 个工作人员组成的池(每个工作人员都将其任务分配给一个代理以计算 doSomething),因此我还需要在此算法中处理效果(例如工作人员故障)。

【问题讨论】:

  • 工作队列 (WQ) 是一个 List[A]。当我使用 doSomething 处理 WQ 中的每个项目时(每个项目返回一个 List[A]),这将以纯粹的功能方式添加到 WQ 的末尾?
  • 感谢@dk14。我也在 akka-streams 和图形 dsl 上阅读了一些内容,也许这是一个竞争者。回顾一下,我们从 A 类型的单个项目开始,该项目通过函数“doSomething”返回项目列表(Nil 或更多 A 项目),然后这些项目中的每一个都反馈到相同的“doSomething”函数相继。也许有一种方法可以在 akka-streams 中使用 doSomething 作为 Flow 函数很好地表达这一点?

标签: scala scalaz akka-stream scalaz-stream


【解决方案1】:

那么“怎么做”的答案呢?是:

import scalaz.stream._
import scalaz.stream.async._
import Process._

def doSomething(i: Int) = if (i == 0) Nil else List(i - 1)

val q = unboundedQueue[Int]
val out = unboundedQueue[Int]

q.dequeue
 .flatMap(e => emitAll(doSomething(e)))
 .observe(out.enqueue)
 .to(q.enqueue).run.runAsync(_ => ()) //runAsync can process failures, there is `.onFailure` as well

q.enqueueAll(List(3,5,7)).run
q.size.continuous
 .filter(0==)
 .map(_ => -1)
 .to(out.enqueue).once.run.runAsync(_ => ()) //call it only after enqueueAll

import scalaz._, Scalaz._
val result = out
  .dequeue
  .takeWhile(_ != -1)
  .map(_.point[List])
  .foldMonoid.runLast.run.get //run synchronously

结果:

result: List[Int] = List(2, 4, 6, 1, 3, 5, 0, 2, 4, 1, 3, 0, 2, 1, 0)

但是,您可能会注意到:

1) 我必须解决终止问题。 akka-stream 也存在同样的问题,而且更难解决,因为您无法访问队列,也没有自然的背压来保证队列不会因为阅读速度快而为空。

2) 我不得不为输出引入另一个队列(并将其转换为 List),因为工作队列在计算结束时变为空。

因此,这两个库都不太适合此类要求(有限流),但是 scalaz-stream(在删除 scalaz 依赖项后将变为“fs2”)足够灵活,可以实现您的想法。最大的“但是”是默认情况下它将按顺序运行。有(至少)两种方法可以让它更快:

1) 将您的 doSomething 拆分为多个阶段,例如 .flatMap(doSomething1).flatMap(doSomething2).map(doSomething3),然后在它们之间放置另一个队列(如果阶段花费相同的时间,大约快 3 倍)。

2) 并行化队列处理。 Akka 有mapAsync - 它可以自动并行执行maps。 Scalaz-stream 有块 - 你可以将你的 q 分成块,比如说 5,然后并行处理块中的每个元素。无论如何,这两种解决方案(akka vs scalaz)都不太适合使用一个队列作为输入和输出。

但是,这又太复杂了而且毫无意义,因为有一个经典的简单方法:

@tailrec def calculate(l: List[Int], acc: List[Int]): List[Int] = 
  if (l.isEmpty) acc else { 
    val processed = l.flatMap(doSomething) 
    calculate(processed, acc ++ processed) 
  }

scala> calculate(List(3,5,7), Nil)
res5: List[Int] = List(2, 4, 6, 1, 3, 5, 0, 2, 4, 1, 3, 0, 2, 1, 0)

这是并行化的:

@tailrec def calculate(l: List[Int], acc: List[Int]): List[Int] = 
  if (l.isEmpty) acc else { 
    val processed = l.par.flatMap(doSomething).toList
    calculate(processed, acc ++ processed) 
  }

scala> calculate(List(3,5,7), Nil)
res6: List[Int] = List(2, 4, 6, 1, 3, 5, 0, 2, 4, 1, 3, 0, 2, 1, 0)

所以,是的,我会说 scalaz-stream 和 akka-streams 都不符合您的要求;但是经典的 scala 并行集合非常适合。

如果您需要跨多个 JVM 的分布式计算 - 看看 Apache Spark,它的 scala-dsl 使用相同的 map/flatMap/fold 样式。它允许您使用不适合 JVM 内存的大型集合(通过跨 JVM 扩展它们),因此您可以通过使用 RDD 而不是 List 来改进 @tailrec def calculate。它还将为您提供处理 doSomething 内部故障的工具。

附:所以这就是为什么我不喜欢使用流媒体库来完成这些任务的原因。流式传输更像是来自某些外部系统(如 HttpRequests)的无限流,而不是预定义(甚至是大)数据的计算。

P.S.2 如果你需要响应式(无阻塞),你可以使用Future(或scalaz.concurrent.Task)+Future.sequence

【讨论】:

  • 感谢 @dk14 抽出宝贵时间思考您所想到的各种解决方案的优点。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-08-18
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多