多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

Akka Streams 快速入门:从第一个 Source 到背压与物化值完整实战指南

Akka Streams 快速入门:从第一个 Source 到背压与物化值完整实战指南 Akka Streams 快速入门从第一个 Source 到背压与物化值完整实战指南【免费下载链接】akka-coreA platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.项目地址: https://gitcode.com/gh_mirrors/ak/akka-core导读本文基于 Akka 官方快速入门指南stream-quickstart.md展开配套仓库内完整可运行的 Scala/Java 测试示例系统讲解 Akka Streams 的核心用法如何创建并运行一个流、如何复用流蓝图、如何做基于时间的节流处理以及如何用Broadcast分流、用buffer显式处理背压、用物化值Materialized Value获取运行结果。读完本文你将掌握 Akka Streams 从描述到运行的完整心智模型并能够直接上手编写自己的第一个流处理程序。Akka Streams 是 Akka 平台面向弹性Elastic、敏捷Agile、韧性Resilient应用的核心模块之一它基于 Reactive Streams 规范实现了异步、非阻塞、自带背压back-pressure的流式处理。与把数据全部载入内存的集合操作不同Akka Streams 处理的是可以无限流动的数据并且始终以背压机制保护系统不被慢消费者拖垮。1. 环境准备添加依赖要使用 Akka Streams首先在项目中添加akka-stream模块。官方推荐的用法是通过 Akka BOMBill of Materials统一管理版本sbt在build.sbt中声明AkkaVersion变量后使用Maven/Gradle通过akka-bom导入依赖版本管理依赖坐标如下$scala.binary.version$为当前使用的 Scala 二进制版本号例如2.13$akka.version$为本仓库对应的 Akka 版本号// sbt libraryDependencies com.typesafe.akka %% akka-stream % AkkaVersion!-- Maven: 通过 akka-bom 管理版本 -- dependency groupIdcom.typesafe.akka/groupId artifactIdakka-stream_2.13/artifactId /dependency// Gradle: 通过 akka-bom 管理版本 implementation com.typesafe.akka:akka-stream_2.13注意Akka 的依赖托管在 Akka 官方安全库中需要按 https://account.akka.io/token 的说明使用带 token 的安全 URL 才能访问。另外两个实用的提示Java 与 Scala DSL 共用一个 JARakka-stream一个构件同时包含akka.stream.scaladslScala DSL和akka.stream.javadslJava DSL。在使用 IntelliJ / Eclipse 时可关闭 IDE 的自动导入提示避免在 Scala 文件中误导入javadsl或在 Java 文件中误导入scaladsl详见 IDE Tips。本文所有示例均来自仓库中的可运行测试QuickStartDocSpec.scala、QuickStartDocTest.java、TwitterStreamQuickstartDocSpec.scala 与 TwitterStreamQuickstartDocTest.java你可以直接在 IDE 中运行它们。1.1 导入与运行环境流式编程需要导入完整的流式工具集// Scala import akka.stream._ import akka.stream.scaladsl._// Java import akka.stream.*; import akka.stream.javadsl.*;如果要在阅读时直接执行代码示例还需要以下导入// Scala import akka.{ Done, NotUsed } import akka.actor.ActorSystem import akka.util.ByteString import scala.concurrent._ import scala.concurrent.duration._ import java.nio.file.Paths// Java import akka.Done; import akka.NotUsed; import akka.actor.ActorSystem; import akka.util.ByteString; import java.nio.file.Paths; import java.math.BigInteger; import java.time.Duration; import java.util.concurrent.CompletionStage;运行任何流都需要一个ActorSystem——它负责调度流运行所需的线程与资源。在 Scala 中把ActorSystem声明为implicit流就能在运行时自动获取它而不必手动传入在 Java 中则需显式传入// Scala object Main extends App { implicit val system: ActorSystem ActorSystem(QuickStart) // Code here }// Java public class Main { public static void main(String[] argv) { final ActorSystem system ActorSystem.create(QuickStart); // Code here } }完整入口见 QuickStartDocSpec.scala 与 Main.java。2. 第一步创建并运行一个 Source流通常从一个**源Source**开始。下面创建一个发出 1 到 100 整数的简单源// Scala val source: Source[Int, NotUsed] Source(1 to 100)// Java final SourceInteger, NotUsed source Source.range(1, 100);2.1 理解Source的两个类型参数Source有两个类型参数这是理解 Akka Streams 的关键第一个参数流中元素的类型这里分别是Int和Integer第二个参数物化值materialized value运行这个源时产生的辅助值。例如网络源可以在运行时提供绑定的端口号或对端地址。当不需要产生任何辅助信息时使用akka.NotUsed类型——简单的整数区间就属于这一类运行它只会得到一个NotUsed。创建 Source 只是描述如何发出前 100 个自然数源此刻并未激活。要把这些数字真正取出来必须运行它// Scala source.runForeach(i println(i))// Java source.runForeach(i - System.out.println(i), system);这行代码会为源补上一个消费函数本例中是打印到控制台并把这条小流水线交给一个 Actor 去运行。激活这件事通过方法名中的run体现——Akka Streams 中所有运行流的其他方法都遵循这一命名模式。从源码实现看runForeach本质上是runWith(Sink.foreach(f))而runWith又是toMat(sink)(Keep.right).run()——即把源与一个逐元素执行函数的Sink连接成RunnableGraph后交给Materializer物化运行见 Source.scala 与 Source.scala。所谓 Materializer就是把流的蓝图变成正在运行的流的组件其核心契约定义在 Materializer.scala 中。2.2 优雅地结束runForeach返回的 Future在scala.App或普通 Java 程序中运行上面的代码你会发现程序不会结束——因为ActorSystem从未被终止。幸运的是runForeach返回一个Future[Done]Java 中为CompletionStageDone当流正常完成时它会被解析因此可以在流结束后终止系统// Scala val done: Future[Done] source.runForeach(i println(i)) implicit val ec system.dispatcher done.onComplete(_ system.terminate())// Java final CompletionStageDone done source.runForeach(i - System.out.println(i), system); done.thenRun(() - system.terminate());这一模式在仓库测试中与futureValue/toCompletableFuture().get()配合使用保证测试进程能干净退出见 QuickStartDocSpec.scala。2.3 变换流用scan计算阶乘并写入文件Akka Streams 中Source只是一份建筑蓝图blueprint可以被复用并嵌入更大的设计中。比如把整数源变换成阶乘序列再写入文件// Scala val factorials source.scan(BigInt(1))((acc, next) acc * next) val result: Future[IOResult] factorials.map(num ByteString(s$num\n)).runWith(FileIO.toPath(Paths.get(factorials.txt)))// Java final SourceBigInteger, NotUsed factorials source.scan(BigInteger.ONE, (acc, next) - acc.multiply(BigInteger.valueOf(next))); final CompletionStageIOResult result factorials .map(num - ByteString.fromString(num.toString() \n)) .runWith(FileIO.toPath(Paths.get(factorials.txt)), system);这里做了三件事scan对整条流运行累积计算。从BigInt(1)BigInteger.ONE开始依次乘以传入的每个数字scan会先发出初始值再发出每一次的计算结果——于是得到了阶乘数列1!, 2!, 3! …。与fold不同scan把每一步中间结果都作为元素发出适合在这里构建可复用的源。mapByteString把数字序列转换成描述文本文件中每一行的ByteString对象。runWith(FileIO.toPath(...))把流接到文件这个接收端上运行。在 Akka Streams 术语中这种数据接收端叫做Sink。IOResult是 Akka Streams 中 IO 操作返回的类型用于告知处理了多少字节/元素以及流是正常结束还是异常结束。它的结构定义在 IOResult.scalacount表示处理的字节/元素数量status在 2.6.0 之后恒为Success(Done)。再次强调执行到这一步时什么都没有真正计算。factorials只是对将来要计算什么的一份描述只有调用runWith/run之类的运行方法后计算才会发生。3. 可复用组件Flow、Sink 与 Keep.rightAkka Streams 的一个突出特性也是其他流库通常不具备的是不仅 Source 可以像蓝图一样复用所有元素都可以。我们可以把文件写入 Sink与把字符串转成ByteString的处理步骤打包成一个可复用组件。由于流的写法总是从左到右就像自然语言需要一个像 Source 一样但输入口开放的起点——这就是Flow// Scala def lineSink(filename: String): Sink[String, Future[IOResult]] Flow[String].map(s ByteString(s \n)).toMat(FileIO.toPath(Paths.get(filename)))(Keep.right)// Java public SinkString, CompletionStageIOResult lineSink(String filename) { return Flow.of(String.class) .map(s - ByteString.fromString(s.toString() \n)) .toMat(FileIO.toPath(Paths.get(filename)), Keep.right()); }从字符串流出发把每个字符串转换为ByteString再喂给已知的文件写入 Sink得到的蓝图类型是Sink[String, Future[IOResult]]JavaSinkString, CompletionStageIOResult——它接受字符串作为输入物化时会生成Future[IOResult]类型的辅助信息。3.1 为什么需要Keep.right在Source或Flow上链式操作时辅助信息物化值的类型由最左侧起点决定。这里我们想保留FileIO.toPath这个 Sink 提供的物化值Future[IOResult]而不是Flow[String].map(...)的物化值NotUsed所以需要用Keep.right明确告诉实现只关心右侧最新接上算子的物化类型。与之相对的是Keep.left保留左侧与Keep.both两者都保留。现在把这个新 Sink 接到factorials源上稍作适配把数字转成字符串// Scala factorials.map(_.toString).runWith(lineSink(factorial2.txt))// Java factorials.map(BigInteger::toString).runWith(lineSink(factorial2.txt), system);lineSink的完整定义可查看 QuickStartDocSpec.scala 与 QuickStartDocTest.java。4. 基于时间的处理zipWith 与 throttle前面所有操作都与时间无关在严格的元素集合上同样可以完成。下面这个例子展示了流的真正流动性用zipWith把factorials流与另一个发出 0 到 100 的Source拉链合并拼成3! 6这样的字符串// Scala factorials .zipWith(Source(0 to 100))((num, idx) s$idx! $num) .throttle(1, 1.second) .runForeach(println)// Java factorials .zipWith(Source.range(0, 99), (num, idx) - String.format(%d! %s, idx, num)) .throttle(1, Duration.ofSeconds(1)) .runForeach(s - System.out.println(s), system);zipWith按位置两两配对factorials发出的第一个元素是 0 的阶乘第二个是 1 的阶乘依此类推throttle(1, 1.second)Java 用Duration把流减速到每秒 1 个元素。运行该程序你会看到每隔一秒打印一行。4.1 背压Back-pressure初体验有一个不易察觉却至关重要的点如果把两股流各放大到发出十亿个数字JVM 也不会因OutOfMemoryError崩溃——尽管流是在后台异步运行的这正是辅助信息以Future/CompletionStage形式提供的原因。背后的秘密是Akka Streams 隐式地实现了贯穿全链路的流控所有算子都尊重背压。当throttle只能接受每秒 1 个元素时它会向上游的所有数据源发出信号声明自己只能以特定速率接受元素——当上游速率超过每秒 1 个时throttle算子就会对上游施加背压。关于背压协议Reactive Streams 兼容实现共享的详细原理见 Back-pressure explained。5. Reactive Tweets一个完整的流处理案例流处理的一个典型场景是消费一条实时的数据流从中提取或聚合信息。下面以从推文流中提取与 Akka 相关的信息为例展开。这里还涉及所有非阻塞流式方案都要面对的根本问题如果订阅者消费实时数据流太慢怎么办传统方案往往是缓冲元素但缓冲最终通常会导致溢出和系统不稳定。Akka Streams 的答案则是依赖内部背压信号来精确控制这种场景下应该发生什么。5.1 数据模型快速入门示例贯穿全程的数据模型Scala 版完整定义见 TwitterStreamQuickstartDocSpec.scalaJava 版见 TwitterStreamQuickstartDocTest.java// Scala final case class Author(handle: String) final case class Hashtag(name: String) final case class Tweet(author: Author, timestamp: Long, body: String) { def hashtags: Set[Hashtag] body .split( ) .collect { case t if t.startsWith(#) Hashtag(t.replaceAll([^#\\w], )) } .toSet } val akkaTag Hashtag(#akka)提示如果你想先了解术语概览再进入实例可以阅读 Core concepts 与 Defining and running streams 两节再回到本文把各知识点拼装成一个完整应用。5.2 环境与推文源准备环境创建负责运行流的ActorSystem// Scala implicit val system: ActorSystem ActorSystem(reactive-tweets)// Java final ActorSystem system ActorSystem.create(reactive-tweets);假设我们已经有一条现成的推文流在 Akka 中它被表达为Source[Tweet, NotUsed]JavaSourceTweet, NotUsed// Scala val tweets: Source[Tweet, NotUsed]// Java SourceTweet, NotUsed tweets;流总是从Source[Out, M1]出发穿过一个或多个Flow[In, Out, M2]算子最终被Sink[In, M3]消费。第一个类型参数这里是Tweet表示源产生的元素类型M系列类型参数描述的是物化期间创建的对象详见下文物化值小节——NotUsed表示不产生任何值相当于泛型意义上的void。5.3 变换与消费简单流filter / map / runWith下面的操作对熟悉 Scala Collections 的人来说似曾相识但它们作用在流上而非数据集合上这是一个非常重要的区别有些操作只有在流式语境下才有意义反之亦然// Scala val authors: Source[Author, NotUsed] tweets.filter(_.hashtags.contains(akkaTag)).map(_.author)// Java final SourceAuthor, NotUsed authors tweets.filter(t - t.hashtags().contains(AKKA)).map(t - t.author);要把这条流水线物化并运行需要把 Flow 接到一个能让它跑起来的Sink上。最简单的做法是在Source上调用runWith(sink)。为了方便Sink伴生对象Java 中为Sink类的静态方法预定义了许多常用 Sink。现在先把每位作者打印出来// Scala authors.runWith(Sink.foreach(println))// Java authors.runWith(Sink.foreach(a - System.out.println(a)), system);或者使用快捷写法只为最流行的 Sink 如Sink.fold、Sink.foreach定义// Scala authors.runForeach(println)// Java authors.runForeach(a - System.out.println(a), system);完整的首个示例代码见 TwitterStreamQuickstartDocSpec.scala 与 TwitterStreamQuickstartDocTest.java。物化和运行流始终需要一个ActorSystem——在 Scala 中置于隐式作用域或显式传入在 Java 中显式传入如runWith(sink, system)或runWith(sink, materializer)。5.4 流式扁平化mapConcat前面处理的是元素间 1:1 的关系最常见的场景但有时我们需要把一个元素映射成多个元素得到一条扁平化的流类似 Scala Collections 的flatMap。要从推文流中得到扁平化的话题标签流可使用mapConcat算子// Scala val hashtags: Source[Hashtag, NotUsed] tweets.mapConcat(_.hashtags.toList)// Java final SourceHashtag, NotUsed hashtags tweets.mapConcat(t - new ArrayListHashtag(t.hashtags()));关于命名flatMap被刻意回避的原因官方文档给出了两点说明在有界流处理中按连接concatenation方式扁平化往往并不可取有死锁风险merge是更优策略由于活性liveness问题flatMap 的 monad 定律在我们这个实现中并不成立。另外请注意mapConcat要求提供的函数返回一个严格集合Scala 为f: Out immutable.Iterable[T]Java 为Out - java.util.ListT而flatMap则必须全程操作流。6. 广播一条流GraphDSL 与 Broadcast现在假设我们要持久化一条实时流中的全部话题标签和全部作者名——例如把所有作者 handle 写入一个文件、把所有话题标签写入另一个文件。这意味着要把源流拆分成两条流分别处理。在 Akka Streams 中用于构成这种扇出fan-out/ 扇入fan-in结构的元素被称为junction连接点。本例使用Broadcast它把输入端口收到的每个元素原样发给所有输出端口。Akka Streams 刻意将线性的流结构Flows与非线性、分支的结构Graphs分开目的是为两种场景分别提供最顺手的 API。Graph 可以表达任意复杂的流拓扑代价是不如集合变换那样读起来直观。Graph 使用GraphDSL构建// Scala val g RunnableGraph.fromGraph(GraphDSL.create() { implicit b import GraphDSL.Implicits._ val bcast b.add(BroadcastTweet) tweets ~ bcast.in bcast.out(0) ~ Flow[Tweet].map(_.author) ~ writeAuthors bcast.out(1) ~ Flow[Tweet].mapConcat(_.hashtags.toList) ~ writeHashtags ClosedShape }) g.run()// Java RunnableGraph.fromGraph( GraphDSL.create( b - { final UniformFanOutShapeTweet, Tweet bcast b.add(Broadcast.create(2)); final FlowShapeTweet, Author toAuthor b.add(Flow.of(Tweet.class).map(t - t.author)); final FlowShapeTweet, Hashtag toTags b.add( Flow.of(Tweet.class) .mapConcat(t - new ArrayListHashtag(t.hashtags()))); final SinkShapeAuthor authors b.add(writeAuthors); final SinkShapeHashtag hashtags b.add(writeHashtags); b.from(b.add(tweets)).viaFanOut(bcast).via(toAuthor).to(authors); b.from(bcast).via(toTags).to(hashtags); return ClosedShape.getInstance(); })) .run(system);可以看到Scala 版在GraphDSL内部借助隐式图构建器b用~边操作符也可读作 connect / via / to可变地构建图——它由import GraphDSL.Implicits._提供Java 版则用图构建器b组合UniformFanOutShape与Flow。6.1 Graph 与 RunnableGraphGraphDSL.create返回一个Graph本例是Graph[ClosedShape, NotUsed]JavaGraphClosedShape,NotUsed。ClosedShape表示一个完全连通的图closed——没有未连接的输入或输出。正因为它是封闭的才能用RunnableGraph.fromGraph把它转换为RunnableGraph进而调用run()物化出一条正在运行的流。Graph与RunnableGraph都是不可变、线程安全、可自由共享的。图还可以有其他形状——带一个或多个未连接端口表达的是部分图partial graph。关于图的组合与嵌套详见 Modularity, Composition and Hierarchy把复杂计算图包装成 Flow / Sink / Source 的方法详见 Constructing Sources, Sinks and Flows from Partial GraphsJavaConstructing and combining Partial Graphs。7. 背压实战buffer 与 OverflowStrategyAkka Streams 的一大优势是始终把背压信息从流的 Sink订阅方传播到 Source发布方。它不是可选特性而是时刻开启的。不用 Akka Streams 的应用常面临这样的典型问题处理数据的速度暂时或设计上跟不上输入于是开始缓冲输入数据直到没有空间缓冲最终导致OutOfMemoryError或服务响应能力严重下降。在 Akka Streams 中缓冲必须且可以显式处理。例如只要最近 10 条推文可以这样表达// Scala tweets.buffer(10, OverflowStrategy.dropHead).map(slowComputation).runWith(Sink.ignore)// Java tweets .buffer(10, OverflowStrategy.dropHead()) .map(t - slowComputation(t)) .runWith(Sink.ignore(), system);slowComputation模拟每次 500ms 的重计算见 TwitterStreamQuickstartDocSpec.scala。buffer算子需要一个显式且必填的OverflowStrategy它定义缓冲已满又收到新元素时缓冲如何反应。从 OverflowStrategy.scala 源码可以看到完整策略族策略行为是否施加背压dropHead丢弃最旧的元素否dropTail丢弃最新的元素否dropBuffer丢弃整个缓冲否dropNew丢弃新来的元素否backpressure缓冲满时对上游施加背压是fail缓冲满时让流失败错误级别日志否请务必根据实际用例挑选最合适的策略例如dropHead适合只关心最新数据的场景而backpressure适合要求不丢数据的场景。7.1 一个直观的背压演示仓库测试中还提供了一个交互式演示#backpressure-by-readline让流每次发出元素后都等待用户按下回车才继续从而直观感受背压如何沿链路向上游传播// Scala val completion: Future[Done] Source(1 to 10).map(i { println(smap $i); i }).runForeach { i readLine(sElement $i; continue reading? [press enter]\n) }// Java final CompletionStageDone completion Source.from(Arrays.asList(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)) .map(i - { System.out.println(map i); return i; }) .runForeach( i - System.console().readLine(Element %s continue reading? [press enter]\n, i), mat);8. 物化值Materialized Values从流中拿到运行结果到目前为止我们只是用 Flow 处理数据并消费到某种外部 Sink打印或存入外部系统。但有时我们关心从物化后的流水线中获得某个值——例如我们一共处理了多少条推文。先看一个用Sink.fold实现的元素计数器Java 版还需配合Flow.of(Class)使用// Scala val count: Flow[Tweet, Int, NotUsed] Flow[Tweet].map(_ 1) val sumSink: Sink[Int, Future[Int]] Sink.foldInt, Int(_ _) val counterGraph: RunnableGraph[Future[Int]] tweets.via(count).toMat(sumSink)(Keep.right) val sum: Future[Int] counterGraph.run() sum.foreach(c println(sTotal tweets processed: $c))// Java final SinkInteger, CompletionStageInteger sumSink Sink.Integer, Integerfold(0, (acc, elem) - acc elem); final RunnableGraphCompletionStageInteger counter tweets.map(t - 1).toMat(sumSink, Keep.right()); final CompletionStageInteger sum counter.run(system); sum.thenAcceptAsync( c - System.out.println(Total tweets processed: c), system.dispatcher());回顾那些神秘的Mat类型参数Source[Out, Mat]、Flow[-In, Out, Mat]、Sink[-In, Mat]Java 版同理——它们代表这些处理部件物化时返回的值的类型。当你把它们链在一起时可以显式组合各自的物化值本例先用Flow[Tweet].map(_ 1)把每条推文变成整数1Java 版直接对tweets源map再用Sink.fold对流的全部Int元素求和结果以Future[Int]JavaCompletionStageInteger提供用toMat(sumSink)(Keep.right)连接——Keep.right这个预定义函数告诉实现只关心当前接在右侧算子的物化类型于是sumSink的物化类型Future[Int]CompletionStageInteger被保留最终RunnableGraph的类型参数也是Future[Int]。这一步尚未物化流水线只是准备了连接好 Sink 的流描述因此它可以被run()——这正是其类型RunnableGraph[Future[Int]]JavaRunnableGraphCompletionStageInteger所指示的。接下来调用run()物化并运行对RunnableGraph[T]调用run()返回类型T本例为Future[Int]CompletionStageInteger完成后包含推文流的总条数如果流失败这个 Future 将以 Failure 结束。8.1 蓝图可多次物化RunnableGraph可以复用并多次物化因为它只是流的蓝图。比如一条消费实时推文流的图上午物化一次、晚上再物化一次两次物化得到的值互不相同// Scala val sumSink Sink.foldInt, Int(_ _) val counterRunnableGraph: RunnableGraph[Future[Int]] tweetsInMinuteFromNow.filter(_.hashtags contains akkaTag).map(t 1).toMat(sumSink)(Keep.right) // 早上物化一次 val morningTweetsCount: Future[Int] counterRunnableGraph.run() // 晚上再次物化复用同一个蓝图 val eveningTweetsCount: Future[Int] counterRunnableGraph.run()// Java final SinkInteger, CompletionStageInteger sumSink Sink.Integer, Integerfold(0, (acc, elem) - acc elem); final RunnableGraphCompletionStageInteger counterRunnableGraph tweetsInMinuteFromNow .filter(t - t.hashtags().contains(AKKA)) .map(t - 1) .toMat(sumSink, Keep.right()); // materialize the stream once in the morning final CompletionStageInteger morningTweetsCount counterRunnableGraph.run(system); // and once in the evening, reusing the blueprint final CompletionStageInteger eveningTweetsCount counterRunnableGraph.run(system);Akka Streams 中的许多元素都能提供物化值用于获取计算结果或操控这些元素本身详见 Stream Materialization。8.2 runWith一行搞定明白了上述机制就能看懂下面这行一行代码它与上面多行版本完全等价// Scala val sum: Future[Int] tweets.map(t 1).runWith(sumSink)// Java final CompletionStageInteger sum tweets.map(t - 1).runWith(sumSink, system);runWith()是一个便捷方法它自动忽略除runWith()自身所接算子之外所有算子的物化值。在上例中它等价于用Keep.right作为物化值的组合器。9. 总结与后续学习路径通过本快速入门你已经掌握了 Akka Streams 的核心心智模型描述与运行分离Source/Flow/Sink只是蓝图只有run()/runWith()/runForeach()等运行方法才会真正物化并启动流物化值Source[Out, Mat]的第二类型参数让你能拿到运行结果如Future[IOResult]、Future[Int]配合Keep.right/Keep.left组合器精确控制保留哪一个可复用性Source、Flow、Sink、RunnableGraph 都是不可变、线程安全、可自由共享、可多次物化的蓝图背压贯穿始终背压时刻开启throttle、buffer等算子显式表达速率与缓冲策略OverflowStrategy六种策略任选线性与非线性统一线性结构用 Flow 链式 API分支/汇聚结构用GraphDSLBroadcast等 junction 表达。想进一步深入可以继续阅读仓库文档Streams 操作符索引涵盖数十种 Source/Sink 与更多变换算子、Streams 基础与核心概念、流组合与模块化 以及 流图Graphs。所有代码示例均可在仓库的akka-docs/src/test目录下找到并直接运行从第一个Source(1 to 100)到完整的 Reactive Tweets 案例一步步搭建你自己的流处理应用。【免费下载链接】akka-coreA platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.项目地址: https://gitcode.com/gh_mirrors/ak/akka-core创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表