
后端并发编程异步编程【免费下载链接】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 官方文档《Basics and working with Flows》为主体系统讲解 Akka Streams 的核心抽象Source、Sink、Flow、RunnableGraph、有界缓冲boundedness与异步非阻塞背压协议、流的物化materialization机制以及算子融合Operator Fusion、异步边界async boundary、物化值组合、Source 预物化、流的有序性保证与 Actor Materializer 生命周期管理等实战要点。读完本文你将能够正确搭建 akka-stream 依赖、构建并运行线性处理流水线理解背压在慢生产者/快消费者与快生产者/慢消费者两种场景下的行为差异并掌握在多 Actor 场景下绑定或解绑流生命周期的正确姿势。1. 依赖引入Akka Streams 是 Akka 的核心模块之一在使用前需要先在构建工具中加入akka-stream依赖。官方文档推荐通过 Akka BOMakka-bom_$scala.binary.version$统一管理版本避免手工维护多个模块的版本号。以 sbt 为例在build.sbt中声明val AkkaVersion 2.9.x // 以仓库当前版本为准 libraryDependencies com.typesafe.akka %% akka-stream % AkkaVersionMaven 方式pom.xmlproperties akka.version2.9.x/akka.version scala.binary.version2.13/scala.binary.version /properties dependencyManagement dependencies dependency groupIdcom.typesafe.akka/groupId artifactIdakka-bom_${scala.binary.version}/artifactId version${akka.version}/version typepom/type scopeimport/scope /dependency /dependencies /dependencyManagement dependency groupIdcom.typesafe.akka/groupId artifactIdakka-stream_${scala.binary.version}/artifactId /dependencyGradle 方式类似通过 BOM 导入后直接声明com.typesafe.akka:akka-stream_$scala.binary.version$即可。仓库内的 BOM 定义可参考 artifact-bom。2. 核心概念有界缓冲与五要素Akka Streams 是一个使用有界缓冲空间处理和传输元素序列的库。所谓有界性boundedness是其区别于 Actor 模型的关键特性流水线中的每个处理实体独立可能并发执行任意时刻只缓冲有限数量的元素。这与 Actor 邮箱通常无界或有界但会丢弃消息不同——流处理实体的邮箱是有界的且不会丢弃消息。文档定义的五个基础术语贯穿整个文档体系术语含义Stream一个活跃的、涉及数据移动与转换的过程Element流的处理单元所有算子都在上游与下游之间转换、传递元素缓冲大小总是以元素个数为单位与元素实际大小无关Back-pressure一种流控手段消费者将自身当前的可接收能力告知生产者从而有效降低上游生产速率以匹配消费速率在 Akka Streams 中背压始终是非阻塞且异步的Non-Blocking某个操作即使耗时很久也不会阻塞调用线程的进度Graph对流处理拓扑的描述定义元素在流运行时的流动路径Operator构成 Graph 的所有构建块的统称例如map()、filter()、自定义的GraphStage以及Merge、Broadcast等图连接件完整内置算子清单见 算子索引所谓异步、非阻塞背压指的是 Akka Streams 的算子之间通过异步消息传递交换数据而非阻塞调用因此可以减慢快速生产者而不会阻塞其线程——等待中的实体等待慢消费者的快生产者不会霸占线程而是把线程归还给底层线程池这是对线程池友好的设计。3. 定义与运行流四大抽象线性处理流水线由四个核心抽象构成Source恰好一个输出的算子当下游就绪时发射数据元素Sink恰好一个输入的算子请求并接收数据元素可能会拖慢上游生产者Flow恰好一个输入和一个输出的算子通过转换流经它的元素来连接上下游RunnableGraph两端分别接上了 Source 和 Sink、随时可以run()的 Flow。可以把Flow附加到Source上得到复合 Source也可以把Flow前置到Sink上得到新的 Sink。当流的两端都接好之后它就表现为RunnableGraph类型——意味着可以执行了。3.1 物化Materialization从蓝图到运行即使构建完 RunnableGraph在物化之前不会有任何数据流动。物化是为 Graph 描述的计算分配全部运行所需资源的过程在 Akka Streams 中通常意味着启动支撑处理的 Actor也可能是打开文件、Socket 连接等取决于流的需要。关键特性Flow 是流水线的描述因此不可变、线程安全、可自由共享——例如可以安全地在 Actor 之间共享或发送让一个 Actor 准备任务、在代码中完全不同的位置物化执行。Scala 下的分步物化示例完整可运行代码见 FlowDocSpec.scalaval source Source(1 to 10) val sink Sink.foldInt, Int(_ _) // 连接 Source 与 Sink得到 RunnableGraph val runnable: RunnableGraph[Future[Int]] source.toMat(sink)(Keep.right) // 物化流取得 Sink 的物化值 val sum: Future[Int] runnable.run()Java 对应版本见 FlowDocTest.javafinal SourceInteger, NotUsed source Source.from(Arrays.asList(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)); final SinkInteger, CompletionStageInteger sink Sink.fold(0, Integer::sum); final RunnableGraphCompletionStageInteger runnable source.toMat(sink, Keep.right()); final CompletionStageInteger sum runnable.run(system);3.2 物化值Mat 与 Keep物化RunnableGraph[T]后Scala 端会得到类型为 T 的物化值materialized value。每个流算子都能产生一个物化值由用户负责把它们组合成新类型。上例中toMat表示要转换 Source 与 Sink 的物化值Keep.right是便捷函数表示只关心 Sink 的物化值。Sink.fold物化出的Future代表整个流上的折叠结果。Java 端在物化RunnableGraph后得到的特殊容器对象称为MaterializedMapSource 与 Sink 是否向其中放入对象由实现决定——例如Sink.fold会在其中放一个代表折叠结果的CompletionStage。由于流可以物化多次每次物化的物化值都会重新计算通常每次返回不同的值。下面的例子中同一个runnable被物化两次两次得到的是不同的FutureScalaval sink Sink.foldInt, Int(_ _) val runnable: RunnableGraph[Future[Int]] Source(1 to 10).toMat(sink)(Keep.right) val sum1: Future[Int] runnable.run() val sum2: Future[Int] runnable.run() // sum1 和 sum2 是不同 Future3.3 runWith一步到位流可能暴露多个物化值但常见需求只关心 Source 或 Sink 的值。为此提供了便捷方法runWith()Sink的runWith需要一个SourceSource的runWith需要一个SinkFlow的runWith需要同时给定Source和Sink因为 Flow 两端都未连接。val source Source(1 to 10) val sink Sink.foldInt, Int(_ _) // 物化流直接拿到 Sink 的物化值 val sum: Future[Int] source.runWith(sink)3.4 算子的不可变性由于算子不可变连接它们会返回新算子而不是修改现有实例——构建长 Flow 时记得把新值赋给变量或直接运行。以下source.map(_ 0)对原source没有任何影响因为map返回的是新SourceScalaval source Source(1 to 10) source.map(_ 0) // 对 source 无影响因为它不可变 source.runWith(Sink.fold(0)(_ _)) // 55 val zeroes source.map(_ 0) // 返回带 map 的新 Source[Int] zeroes.runWith(Sink.fold(0)(_ _)) // 0注意默认情况下 Akka Streams 元素只支持一个下游算子。把扇出fan-out做成显式 opt-in 特性能让默认流元素更简单高效同时通过broadcast向所有下游发信号或balance向某个可用下游发信号等具名扇出元素对多播场景的精确处理方式保有更大的灵活性。3.5 定义 Source、Sink 与 FlowSource与Sink对象提供了丰富的构建方式以下是文档与 FlowDocSpec.scala 中最常用的构造// 从 Iterable 创建 Source Source(List(1, 2, 3)) // 从 Future 创建 Source Source.future(Future.successful(Hello Streams!)) // 从单个元素创建 Source Source.single(only one element) // 空 Source Source.empty // 折叠整个流、以最终结果 Future 作为物化值的 Sink Sink.foldInt, Int(_ _) // 以流的第一个元素Future 作为物化值的 Sink Sink.head // 消费流但不做任何事的 Sink Sink.ignore // 对每个元素执行副作用调用的 Sink Sink.foreachString)3.6 多种接线方式文档给出了多种连接 Source、Sink、Flow 的方式Scala// 显式创建并接线 Source、Sink 和 Flow Source(1 to 6).via(Flow[Int].map(_ * 2)).to(Sink.foreach(println(_))) // 从 Source 出发 val source Source(1 to 6).map(_ * 2) source.to(Sink.foreach(println(_))) // 从 Sink 出发 val sink: Sink[Int, NotUsed] Flow[Int].map(_ * 2).to(Sink.foreach(println(_))) Source(1 to 6).to(sink) // 内联广播到一个 Sink val otherSink: Sink[Int, NotUsed] Flow[Int].alsoTo(Sink.foreach(println(_))).to(Sink.ignore) Source(1 to 6).to(otherSink)Java 侧Source.from(...)、Flow.of(Integer.class).map(...)、source.to(sink)的写法与 Scala 一一对应完整示例见 FlowDocTest.java。3.7 非法流元素null 禁令依据 Reactive Streams 规范Rule 2.13Akka Streams 不允许null作为元素在流中传递。若需要建模值缺失的概念Scala 推荐用Option或EitherJava 推荐用java.util.Optional。4. 背压原理解析Akka Streams 实现了 Reactive Streams 规范定义的异步非阻塞背压协议Akka 是该规范的创始成员之一。库的使用者无需编写任何显式背压处理代码——所有内置算子都自动内置并处理背压。当然也可以通过带溢出策略的显式buffer算子影响流的行为这在包含环路的复杂图中尤其重要环路必须极其小心见 Graph 环、活性与死锁。背压协议以下游Subscriber能接收并缓冲的元素个数来定义这个数量称为demand需求。数据源Reactive Streams 术语中的PublisherAkka Streams 中实现为Source保证绝不向任何Subscriber发射超过其已接收总需求数量的元素。注意Reactive Streams 规范以Publisher/Subscriber定义协议但这两个类型不是面向用户的 API而是不同 Reactive Streams 实现之间的底层构件。Akka Streams 将其实现为Source、Flow对应规范中的Processor、Sink不直接暴露 Reactive Streams 接口。需要与其他响应式流库集成时参见 与 Reactive Streams 集成。这种背压工作模式可以通俗地称为动态 push / pull 模式根据下游能否跟上上游生产速率在基于 push 与基于 pull 的背压模型之间切换。4.1 慢 Publisher、快 Subscriberpush 模式这是理想情况——无需放慢 Publisher。但信号速率很少恒定可能随时变化突然变成Subscriber 比 Publisher 慢因此背压协议在这种场景下也必须保持启用同时不希望为这个安全网付出过高代价。协议通过 Subscriber 异步向 Publisher 发送Request(n)信号解决协议保证 Publisher 发射的元素不会超过已声明的 demand。由于当前 Subscriber 更快它会以更高频率发送 Request 信号也可能批量合并 demand一次请求多个元素。这意味着 Publisher 发射传入元素时几乎不需要等待被背压。此场景实际运行在push 模式Publisher 能多快就多快地产出因为待处理的 demand 会在发射元素时被恰好及时地补充。4.2 快 Publisher、慢 Subscriberpull 模式此时必须对 Publisher 施加背压。由于 Publisher 不允许发射超过 Subscriber 已声明 demand 的元素它只能采用以下策略之一如果能够控制生产速率就不生成元素以有界方式缓冲元素直到收到更多 demand丢弃元素直到收到更多 demand如果以上策略都无法实施就拆除流。此场景实际意味着 Subscriber 从 Publisher 处pull元素称为基于 pull 的背压。5. 流的物化机制在 Akka Streams 中构建 Flow 与图时可以把它们理解为准备蓝图、执行计划。物化Stream Materialization就是拿流描述RunnableGraph并分配其运行所需全部资源的过程——通常意味着启动驱动处理的 Actor也可能打开文件、Socket 等。物化由所谓的终结操作触发最主要的是定义在Source/Flow上的各种run()、runWith()以及少量针对知名 Sink 的语法糖例如runForeach(el ...)即runWith(Sink.foreach(el ...))的别名。物化由 ActorSystem 全局的Materializer在物化线程上同步执行真正的流处理由物化期间启动的 Actor 完成运行在它们被配置到的线程池上。默认线程池是ActorSystem配置中设置的 dispatcher但可以通过给以下两者提供Attributes来指定其他线程池待物化的流以defaultAttributes创建的自定义Materializer实例。注意在复合图中复用线性算子Source、Sink、Flow的实例是合法的但该算子会被物化多次。5.1 算子融合Operator Fusion默认情况下Akka Streams 会融合流算子——一个 Flow 或流的多个处理步骤可在同一个 Actor 内执行带来两个后果融合算子之间传递元素快得多省去了异步消息开销融合的流算子不会彼此并行每个融合部分最多只用一个 CPU 核。要并行处理就必须手动插入异步边界通过Attributes.asyncBoundary即 Source、Sink、Flow 上的async方法把算子标记为以异步方式与其下游通信。文档示例ScalaSource(List(1, 2, 3)).map(_ 1).async.map(_ * 2).to(Sink.ignore)这个例子在 Flow 内创建了两个区域各自在一个 Actor 内执行——如果加一和乘二是极其昂贵的操作两个 CPU 可并行处理从而获得性能提升。异步边界不是流中元素异步传递的单点不像其他流库那样而是始终以向已构建的流图累加信息的方式工作即红泡内的所有内容由一个 Actor 执行红泡外由另一个执行。该方案可连续套用每个边界总是包住前一个边界以及之后新增的全部算子。警告在未融合时代2.0-M2 之前每个流算子都有一个隐式输入缓冲以提升效率。如果流图包含环这些缓冲可能对避免死锁至关重要。融合之后这些隐式缓冲不复存在融合算子之间无缓冲传递数据。在必须缓冲流才能运行的场景需要用.buffer()算子显式插入缓冲——通常大小为 2 的缓冲就足以让反馈环工作。5.2 组合物化值既然每个算子物化后都能提供一个物化值就必须表达把这些算子插接在一起时如何组合成最终值。为此许多算子方法都提供了带额外组合函数参数的变体。以下是 FlowDocSpec.scala 中的经典组合示例// 可从外部显式发信号的 Source val source: Source[Int, Promise[Option[Int]]] Source.maybe[Int] // 内部以 1 个/秒节流、返回 Cancellable 的 Flow val flow: Flow[Int, Int, Cancellable] throttler // 以返回的 Future 携带流中第一个元素的 Sink val sink: Sink[Int, Future[Int]] Sink.head[Int] // 默认保留最左侧算子的物化值 val r1: RunnableGraph[Promise[Option[Int]]] source.via(flow).to(sink) // 用 Keep.right 简单选择物化值 val r2: RunnableGraph[Cancellable] source.viaMat(flow)(Keep.right).to(sink) val r3: RunnableGraph[Future[Int]] source.via(flow).toMat(sink)(Keep.right) // runWith 总是给出 runWith 自身所加算子的物化值 val r4: Future[Int] source.via(flow).runWith(sink) val r5: Promise[Option[Int]] flow.to(sink).runWith(source) val r6: (Promise[Option[Int]], Future[Int]) flow.runWith(source, sink) // 更复杂的组合 val r7: RunnableGraph[(Promise[Option[Int]], Cancellable)] source.viaMat(flow)(Keep.both).to(sink) val r9: RunnableGraph[((Promise[Option[Int]], Cancellable), Future[Int])] source.viaMat(flow)(Keep.both).toMat(sink)(Keep.both) // 也可以用 mapMaterializedValue 转换物化值把 r9 的嵌套二元组压平 val r11: RunnableGraph[(Promise[Option[Int]], Cancellable, Future[Int])] r9.mapMaterializedValue { case ((promise, cancellable), future) (promise, cancellable, future) } // 现在可用模式匹配拿到全部物化值 val (promise, cancellable, future) r11.run()关键点Keep.left取左、Keep.right取右、Keep.both取两者组成的二元组mapMaterializedValue可以对物化值做任意变换。Java 侧对应的是Keep.left()、Keep.right()、Keep.both()与mapMaterializedValue(...)参见 FlowDocTest.java。补充在图中也可以从流内部访问物化值详见 在 Graph 内部访问物化值。5.3 Source 预物化preMaterialize有些场景需要在 Source接入图的其他部分之前就拿到它的物化值——这对由物化值驱动的 Source如Source.queue、Source.actorRef、Source.maybe尤其有用。通过 Source 上的preMaterialize()算子可以同时获得它的物化值和另一个 Source后者可用于消费原 Source 的消息。注意它可以被多次物化。Scala 示例文档与测试中的 actorRef 场景val completeWithDone: PartialFunction[Any, CompletionStrategy] { case Done CompletionStrategy.immediately } val matValuePoweredSource Source.actorRefString // 预物化立即拿到 ActorRef 与可用于后续物化的 Source val (actorRef, source) matValuePoweredSource.preMaterialize() actorRef ! Hello! // 把 source 传递到别处进行物化 source.runWith(Sink.foreach(println))从源码看preMaterialize在 Source.scala 中的实现是toMat(Sink.asPublisher(fanout true))(Keep.both).run()即内部通过一个 fanout 的 Reactive StreamsPublisher实现——这意味着会引入一个缓冲且错误不会向上游传播而是变成不带错误细节的取消信号。Java 对应 API 为Source.actorRef(...)配合preMaterialize(system)见 FlowDocTest.java。6. 流的有序性Stream orderingAkka Streams 中几乎所有计算算子都保持输入元素的顺序若输入{IA1,...,IAn}引起输出{OA1,...,OAk}输入{IB1,...,IBm}引起输出{OB1,...,OBl}且所有IAi都先于所有IBi那么OAi也先于OBi。这一性质对mapAsync这类异步算子也成立但存在不保序的版本mapAsyncUnordered它不保持这种顺序。然而处理多输入流的连接件如Merge不保证来自不同输入端口的元素的输出顺序——merge 类操作可能先发射Ai再发射Bi顺序由内部逻辑决定。而Zip这类专门的算子保证输出顺序因为每个输出元素都依赖所有上游元素已被发出——因此 zip 场景的顺序由这一性质定义。如果需要在 fan-in 场景下对元素发射顺序做细粒度控制可以考虑MergePreferred、MergePrioritized或者自定义GraphStage——它给你对合并方式的完全控制权。7. Actor Materializer 生命周期Materializer负责把流蓝图变成运行中的流并产出物化值。一个 ActorSystem 级别的Materializer由 AkkaExtensionSystemMaterializer提供——Scala 通过隐式ActorSystemJava 通过向各种run方法传入ActorSystem因此除非有特殊需求无需关心Materializer的创建。一个可能需要自定义Materializer实例的用例是把 Actor 中物化的所有流绑定到该 Actor 的生命周期Actor 停止或崩溃时流也随之停止。理解Materializer生命周期是与流、Actor 协作的重要一环物化器绑定到它创建时所在的ActorRefFactory的生命周期实际就是ActorSystem或在 Actor 内创建时ActorContext。自 Akka 2.6 起绑定到ActorSystem应改用系统物化器。由系统物化器运行时流会一直运行到ActorSystem关闭若物化器在流运行完成之前关闭流将被突然终止。这与通常的终止方式cancel/complete不同。流的生命周期如此绑定物化器是为了防止泄漏正常操作中不应依赖此机制而应使用KillSwitch或正常的完成信号管理流的生命周期。7.1 绑定到 Actor 生命周期下面的例子在 Actor 内创建Materializer将其生命周期绑定到该 ActorScala见 FlowDocSpec.scalafinal class RunWithMyself extends Actor { implicit val mat: Materializer Materializer(context) Source.maybe.runWith(Sink.onComplete { case Success(done) println(sCompleted: $done) case Failure(ex) println(sFailed: ${ex.getMessage}) }) def receive { case boom context.stop(self) // 也会终止该流 } }这里用ActorContext创建物化器把它的生命周期绑定到外层 Actor正常情况下流会永远运行但如果停止该 Actor流也会被终止——流的生命周期被绑定到了所在 Actor 的生命周期。当 Actor 代表某个实体例如用户且我们用创建的流持续查询该实体时这个技术非常有用——Actor 已终止时还让流存活没有意义。流的终止会以流上的 Abrupt termination exception 信号体现。也可以显式调用Materializer.shutdown()关闭物化器从而突然终止其运行的所有流。7.2 让流超越 Actor 生命周期有时你希望显式创建一个比 Actor 活得更久的流例如用 Akka Stream 向外部服务推送大批数据Actor 已完成全部职责、想尽早停止。此时应把系统物化器传入 Actorfinal class RunForever(implicit val mat: Materializer) extends Actor { Source.maybe.runWith(Sink.onComplete { case Success(done) println(sCompleted: $done) case Failure(ex) println(sFailed: ${ex.getMessage}) }) def receive { case boom context.stop(self) // 不会终止该流它绑定到系统 } }传入物化器后流绑定到整个ActorSystem而非单个 Actor 的生命周期。如果想共享一个物化器或按物化器设置把流分组到特定物化器这也很有用。警告不要在 Actor 内部通过把context.system传给创建逻辑来新建 Actor 物化器这会导致每个这样的 Actor 都创建一个新的Materializer并可能泄漏除非显式关闭。推荐做法是传入现成 Materializer或用 Actor 的context创建。8. 结语与延伸阅读本文覆盖了 Akka Streams 从依赖引入、核心抽象与有界缓冲思想到背压协议两种模式、物化与物化值组合、算子融合与异步边界、Source 预物化、流有序性以及 Actor Materializer 生命周期管理的完整知识链。文中所有代码片段均取自仓库内真实测试文件可直接在 FlowDocSpec.scala 与 FlowDocTest.java 中核对运行。进一步深入可以参考仓库中的相关文档算子索引全部内置算子的速查表自定义 GraphStage自定义算子与精细控制流行为流图与图 DSLfan-in/fan-out 拓扑、图的物化值访问与环的死锁处理与 Reactive Streams 集成跨实现互操作物化器相关实现可阅读 Source.scala 与 SystemMaterializer.scala。赞分享后端并发编程异步编程【免费下载链接】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 Streams flatMapMerge 算子详解有界并发扁平化合并与背压语义Akka Streams flatMapMerge 算子详解有界并发扁平化合并与背压语义 本篇技术指南以 Akka Streams 官方文档 flatMapM后端并发编程异步编程Akka Streams 的 buffer 算子缓冲、背压与溢出策略OverflowStrategy完全指南Akka Streams 的 buffer 算子缓冲、背压与溢出策略OverflowStrategy完全指南 导读 buffer 是 Akka Strea后端并发编程异步编程Akka Streams conflateWithSeed 算子实战背压下的元素聚合与速率解耦Akka Streams conflateWithSeed 算子实战背压下的元素聚合与速率解耦 导读 conflateWithSeed 是 Akka Stre后端并发编程异步编程上一篇Sublime Text编码转换终极指南告别乱码的完整解决方案下一篇小米手表表盘设计终极指南Mi-Create免费工具完全教程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考