【问题标题】:Reproduce akka-stream async output重现 akka-stream 异步输出
【发布时间】:2016-02-07 13:13:36
【问题描述】:

我是akka-stream的新手,所以想问一下如何重现本文中提出的行为http://doc.akka.io/docs/akka-stream-and-http-experimental/2.0.2/scala/stream-rate.html

对于给定的代码

Source(1 to 3)
  .map { i => println(s"A: $i"); i }
  .map { i => println(s"B: $i"); i }
  .map { i => println(s"C: $i"); i }
  .runWith(Sink.ignore)

得到这样的相似

A: 1
A: 2
B: 1
A: 3
B: 2
C: 1
B: 3
C: 2
C: 3

我尝试添加一些随机的Thread.sleep, 从无限迭代器创建一个流。 但是 Akka 根据调试输出总是使用同一个线程进行处理。

所以问题是:如何使用 akka-stream 重现异步行为(每个阶段都应该以异步方式运行)?

【问题讨论】:

  • 异步下是什么意思?您的预期行为是什么?
  • 某事,除了这个 A: 1 B: 1 C: 1 A: 2 B: 2 C: 2 A: 3 B: 3 C: 3
  • 然后在一个地图步骤中完成所有操作?

标签: multithreading scala akka-stream


【解决方案1】:

您看到顺序操作的原因是因为您的所有操作都来自同一个源,因此在同一个异步边界内。要获得您正在寻找的“异步行为”,您需要添加Flows

implicit val actorSystem = ActorSystem()
implicit val actorMaterializer = ActorMaterializer()

Source(1 to 3).via(Flow[Int].map{i => println(s"A: $i"); i })
              .via(Flow[Int].map{i => println(s"B: $i"); i })
              .via(Flow[Int].map{i => println(s"C: $i"); i })
              .runWith(Sink.ignore)

每个 Flow 都会具体化为一个单独的 Actor。注意:要获得真正的并发,ActorSystem 正在运行的线程池必须有 1 个以上的线程。

要记住的一点:ActorSystem 的好处是它承担了对操作的低级别控制的责任,以便开发人员可以专注于“业务逻辑”。这也可能是一个缺点。根据您的 ActorSystem 配置、JVM 配置和硬件配置,操作顺序可能仍然是同步的。

【讨论】:

  • 不保证将如何执行。每个都可以在单独的线程中执行。你也应该记住缓冲
  • @1esha,完全同意。我的建议只是放宽了导致顺序处理的约束之一。我将根据您的评论进行相应更新。谢谢。
猜你喜欢
  • 1970-01-01
  • 2020-07-20
  • 1970-01-01
  • 2016-01-29
  • 1970-01-01
  • 2018-03-08
  • 1970-01-01
  • 2021-09-08
  • 2020-12-07
相关资源
最近更新 更多