Akka Streams是否会利用Akka Actors?

suz*_*omo 4 scala actor akka akka-stream

我开始学习Akka Streams,这是一个使用反压功能处理数据的框架.该图书馆是Akka的一部分,将自己描述为:

Akka是一个工具包和运行时,用于在JVM上构建高度并发,分布式和弹性的消息驱动应用程序.

这些能力来自Akka演员的本质.但是,从我的角度来看,流处理和演员是彼此无关的概念.

问题:Akka Streams是否利用了Akka演员的这些功能?如果是的话,你会解释演员如何帮助流吗?

Qui*_*zie 6

Akka Streams是比演员更高层次的抽象.它是Reactive Streams的一个实现,它构建 actor模型之上.它利用了所有actor的功能,因为它使用了actor.

您甚至可以直接在流的任何部分使用actor.查看ActorPublisher和ActorSubscriber.


Ram*_*gil 5

一个好的起点是akka 流快速入门

是的, anActor用于“具体化” a 的每个 { Source, Flow, Sink} Stream。这意味着当您创建一个 Stream 时,在该流被具体化之前实际上不会发生任何事情,通常是通过.run()方法调用。

例如,这里定义了一个 Stream:

import akka.actor.ActorSystem
import akka.stream.ActorMaterializer
import akka.stream.scaladsl.{Source, Flow, Sink}

val stream = Source.single[String]("test")
                   .via(Flow[String].filter(_.size > 0))
                   .to(Sink.foreach{println})
Run Code Online (Sandbox Code Playgroud)

即使流现在是一个val没有实际发生的计算。Stream 只是一个计算方法。要真正开始工作,需要实现 Stream。下面是一个不使用隐式的例子来清楚地展示物化是如何发生的:

val actorSystem = ActorSystem()
val materializer = ActorMaterializer()(actorSystem)

stream.run()(materializer)  //work begins
Run Code Online (Sandbox Code Playgroud)

现在已经创建了 3 个 Actors(至少):1 个用于Source.single,1 个用于Flow.filter,1 个用于Sink.foreach。注意:您可以使用相同的materializer方法启动其他流

val doesNothingStream = Source.empty[String]
                              .to(Sink.ignore)
                              .run()(materializer)
Run Code Online (Sandbox Code Playgroud)