suz*_*omo 4 scala actor akka akka-stream
我开始学习Akka Streams,这是一个使用反压功能处理数据的框架.该图书馆是Akka的一部分,将自己描述为:
Akka是一个工具包和运行时,用于在JVM上构建高度并发,分布式和弹性的消息驱动应用程序.
这些能力来自Akka演员的本质.但是,从我的角度来看,流处理和演员是彼此无关的概念.
问题:Akka Streams是否利用了Akka演员的这些功能?如果是的话,你会解释演员如何帮助流吗?
Akka Streams是比演员更高层次的抽象.它是Reactive Streams的一个实现,它构建在 actor模型之上.它利用了所有actor的功能,因为它使用了actor.
您甚至可以直接在流的任何部分使用actor.查看ActorPublisher和ActorSubscriber.
一个好的起点是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)
| 归档时间: |
|
| 查看次数: |
520 次 |
| 最近记录: |