【问题标题】:Fold and reduce show non-deterministic behavior when ran in parallel, why?当并行运行时,折叠和减少显示不确定的行为,为什么?
【发布时间】:2020-11-17 17:10:42
【问题描述】:

所以我正在尝试使用 Akka Streams 计算项目的出现次数。 下面的示例是我所拥有的简化版本。我需要两个管道同时工作。由于某种原因,打印的结果不正确。

有人知道为什么会这样吗?我是否遗漏了有关子流的重要内容?

/**
 * SIMPLE EXAMPLE
 */
object TestingObject {
  import akka.actor.ActorSystem
  import akka.stream._
  import akka.stream.scaladsl._
  import java.nio.file.Paths
  import akka.util.ByteString
  import counting._
  import graph_components._

  // implicit actor system
  implicit val system:ActorSystem = ActorSystem("Sys")

  def main(args: Array[String]): Unit = {

    val customFlow = Flow.fromGraph(GraphDSL.create() {
      implicit builder =>
        import GraphDSL.Implicits._


        // Components
        val A   = builder.add(Balance[(Int, Int)](2, waitForAllDownstreams = true));
        val B1   = builder.add(mergeCountFold.async);
        val B2   = builder.add(mergeCountFold.async);
        val C   = builder.add(Merge[(Int, Int)](2));
        val D   = builder.add(mergeCountReduce);

        // Graph
        A ~> B1 ~> C ~> D
        A ~> B2 ~> C

        FlowShape(A.in, D.out);
    })

    // Run
    Source(0 to 101)
      .groupBy(10, x => x % 4)
      .map(x => (x % 4, 1))
      .via(customFlow)
      .mergeSubstreams
      .to(Sink.foreach(println)).run();
  }

  def mergeCountReduce = Flow[(Int, Int)].reduce((l, r) => {
    println("REDUCING");
    (l._1, l._2 + r._2)
  })
  def mergeCountFold = Flow[(Int, Int)].fold[(Int,Int)](0,0)((l, r) => {
    println("FOLDING");
    (r._1, l._2 + r._2)
  })

}

【问题讨论】:

  • 您将mergeCountReduce 的哪些用法替换为mergeCountFold
  • “仅输出部分结果项”和“错过元组的第一个值”究竟是什么意思?请注意,在您的示例中,foldreduce 之间的区别在于,fold 将发出最后一个看到的值的 _1,而reduce 将发出第一个看到的值的_1。由于这些都依赖于订购,还值得注意的是 BalanceMerge 组合(尤其是在两者之间的 async)不提供真正的订购保证。
  • 我正在替换这两种用法,我正在或正在使用 mergeCountReduce 或 mergeCountFold,没有混合。
  • “只有一些值”是指应该打印合并结果的接收器在使用 Reduce 时不会打印所有结果。使用 Fold 时,它会打印所有结果,但其中一些结果为 0 作为 ._1
  • 排序并不重要,因为我之后将所有内容与 mergcountReduce/Fold 合并

标签: scala akka-stream


【解决方案1】:

两个观察结果:

  • mergeCountReduce 将发出它看到的第一个键以及看到的值的总和(如果它没有看到任何元素,它将使流失败)
  • mergeCountFold 将发出它看到的最后一个键和看到的值的总和(如果它没有看到任何元素,它将发出一个键和零值)

(在这两种情况下,虽然密钥始终相同)

这些观察都不受async 边界的影响。

不过,在前面的 Balance 运算符的上下文中,async 引入了一个隐式缓冲区,它可以防止它包裹的图形在缓冲区满之前反压。 Balance 将流值发送到没有背压的第一个输出,因此如果Balance 之后的阶段没有明显慢于上游,Balance 可能只将值发送到一个输出(在这种情况下为B1) .

在这种情况下,reduceB1 将发出密钥和计数,而 B2 失败,导致整个流失败。

对于fold,在这种情况下,B1 会发出键和计数,而B2,没有看到任何值会发出(0,0)。合并将按照它们发出的顺序发出它们(合理地假设有 50/50 的机会),因此最终折叠将要么具有键和计数,要么具有零和计数。

【讨论】:

  • 那么异步是麻烦制造者吗?如果我放弃它,将流分成两部分没有任何优势。有什么解决办法吗?
  • 另外,我不知道你是否已经运行了代码,但是当使用 reduce 版本时,根本不会打印任何内容。知道为什么会这样吗?
  • 其实我都忘了reduce没有元素时会发生什么,正在编辑……
  • 我认为折叠版本也失败了。没有打印任何内容。
  • 抛开当其中一个没有任何元素时会发生什么,我如何确保它们都获得均衡数量的元素?
猜你喜欢
  • 1970-01-01
  • 2020-10-28
  • 2012-07-11
  • 2021-10-01
  • 1970-01-01
  • 2020-03-21
  • 2020-12-12
  • 1970-01-01
  • 2017-10-20
相关资源
最近更新 更多