【问题标题】:Does Akka Streams Unzip/Zip preserve order?Akka Streams Unzip/Zip 是否保留顺序?
【发布时间】:2021-03-02 14:50:25
【问题描述】:

如果我解压缩一系列元组,对两个流执行一些异步突变,然后重新压缩它们,Akka 是否保证流以相同的顺序重新压缩?

例子:

import akka.NotUsed
import akka.actor.ActorSystem
import akka.stream.{ActorMaterializer, FlowShape}
import akka.stream.scaladsl.{Flow, GraphDSL, Sink, Source, Unzip, Zip}

import scala.concurrent.ExecutionContext.Implicits.global
import scala.concurrent.duration.DurationInt
import scala.concurrent.{Await, Future}

val graph: Flow[(Int, String), (Int, String), NotUsed] = Flow.fromGraph(GraphDSL.create() { implicit builder =>
  import GraphDSL.Implicits._

  val unzip = builder.add(Unzip[Int, String])
  val increment = builder.add(Flow[Int].mapAsync(3) { num => Future(num + 1) })
  val append = builder.add(Flow[String].mapAsync(3) { letter => Future(s"$letter-x") })
  val zip = builder.add(Zip[Int, String])

  unzip.out0 ~> increment ~> zip.in0
  unzip.out1 ~> append ~> zip.in1

  FlowShape(unzip.in, zip.out)
})

implicit val system = ActorSystem()
implicit val materializer = ActorMaterializer()

val out = Source(collection.immutable.Seq((0, "a"), (1, "b"), (2, "c")))
  .via(graph)
  .runWith(Sink.seq)

Await.result(out, 1 second)

在这个简单的测试中,输出为Vector((1,a-x), (2,b-x), (3,c-x))。所以事情看起来不错。但我不确定我是否可以相信这将永远如此。

引起一些关注的是:

val unzip = builder.add(Unzip[Int, String])
val increment = builder.add(Flow[Int].mapAsync(3) { num => Future(num + 1) })
val filter = builder.add(Flow[Int].filter(_ != 2))
val append = builder.add(Flow[String].mapAsync(3) { letter => Future(s"$letter-x") })
val zip = builder.add(Zip[Int, String])

unzip.out0 ~> increment ~> filter ~> zip.in0
unzip.out1 ~> append ~> zip.in1

// output: Vector((1,a-x), (3,b-x))

即使保留了顺序,也不能保证会保留原始元组关系。

我可以手动检查我的流程以确保没有过滤逻辑。但是完成后,我可以确定元组会按照它们收到的顺序重新压缩吗?

【问题讨论】:

    标签: scala akka akka-stream


    【解决方案1】:

    TL;DR 是的,确实如此。来自 Akka 中的 Stream ordering 文档:

    在 Akka Streams 中,几乎所有计算运算符都保留元素的输入顺序。这意味着如果输入{IA1,IA2,...,IAn}“原因”输出{OA1,OA2,...,OAk},输入{IB1,IB2,...,IBm}“原因”输出{OB1,OB2,...,OBl},并且所有IAi发生在所有IBi之前,那么OAi发生在OBi之前。

    async 操作(例如 mapAsync)甚至支持此属性,但是存在称为 mapAsyncUnordered 的无序版本,它不保留此顺序。

    但是,对于处理多个输入流的 Junction(例如 Merge),通常不会为到达不同输入端口的元素定义输出顺序。也就是说,类似合并的操作可能会在发出Bi 之前发出Ai,并且由其内部逻辑决定发出元素的顺序。 Zip 等特殊元素确实保证了它们的输出顺序,因为每个输出元素都依赖于所有已发出信号的上游元素——因此在压缩的情况下的顺序由该属性定义。

    如果您发现自己需要在扇入场景中对发射元素的顺序进行细粒度控制,请考虑使用MergePreferredMergePrioritizedGraphStage——这使您可以完全控制合并的执行方式。

    【讨论】:

    • 甜蜜。我还发现了一个支持这一结论的PassThrough flow 示例。谢谢!
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2010-10-14
    • 2017-03-04
    • 1970-01-01
    • 2021-03-04
    • 2016-02-12
    • 1970-01-01
    相关资源
    最近更新 更多