Akka Streams 基础与 Flow 实战:从有界缓冲、背压到物化与算子融合
2026/9/24 16:06:13 网站建设 项目流程
  • 后端
  • 并发编程
  • 异步编程

【免费下载链接】akka-core

A 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 BOM(akka-bom_$scala.binary.version$)统一管理版本,避免手工维护多个模块的版本号。

以 sbt 为例,在build.sbt中声明:

val AkkaVersion = "2.9.x" // 以仓库当前版本为准 libraryDependencies += "com.typesafe.akka" %% "akka-stream" % AkkaVersion

Maven 方式(pom.xml):

<properties> <akka.version>2.9.x</akka.version> <scala.binary.version>2.13</scala.binary.version> </properties> <dependencyManagement> <dependencies> <dependency> <groupId>com.typesafe.akka</groupId> <artifactId>akka-bom_${scala.binary.version}</artifactId> <version>${akka.version}</version> <type>pom</type> <scope>import</scope> </dependency> </dependencies> </dependencyManagement> <dependency> <groupId>com.typesafe.akka</groupId> <artifactId>akka-stream_${scala.binary.version}</artifactId> </dependency>

Gradle 方式类似,通过 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,以及MergeBroadcast等图连接件;完整内置算子清单见 算子索引

所谓"异步、非阻塞背压",指的是 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.scala):

val 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.java:

final Source<Integer, NotUsed> source = Source.from(Arrays.asList(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)); final Sink<Integer, CompletionStage<Integer>> sink = Sink.fold(0, Integer::sum); final RunnableGraph<CompletionStage<Integer>> runnable = source.toMat(sink, Keep.right()); final CompletionStage<Integer> 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后得到的特殊容器对象称为MaterializedMap,Source 与 Sink 是否向其中放入对象由实现决定——例如Sink.fold会在其中放一个代表折叠结果的CompletionStage

由于流可以物化多次,每次物化的物化值都会重新计算,通常每次返回不同的值。下面的例子中,同一个runnable被物化两次,两次得到的是不同的Future(Scala):

val 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 是不同 Future!

3.3 runWith:一步到位

流可能暴露多个物化值,但常见需求只关心 Source 或 Sink 的值。为此提供了便捷方法runWith()SinkrunWith需要一个SourceSourcerunWith需要一个SinkFlowrunWith需要同时给定SourceSink(因为 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返回的是新Source(Scala):

val 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 与 Flow

SourceSink对象提供了丰富的构建方式,以下是文档与 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.13),Akka Streams 不允许null作为元素在流中传递。若需要建模"值缺失"的概念,Scala 推荐用OptionEither,Java 推荐用java.util.Optional

4. 背压原理解析

Akka Streams 实现了 Reactive Streams 规范定义的异步非阻塞背压协议(Akka 是该规范的创始成员之一)。库的使用者无需编写任何显式背压处理代码——所有内置算子都自动内置并处理背压。当然,也可以通过带溢出策略的显式buffer算子影响流的行为,这在包含环路的复杂图中尤其重要(环路必须极其小心,见 Graph 环、活性与死锁)。

背压协议以"下游Subscriber能接收并缓冲的元素个数"来定义,这个数量称为demand(需求)。数据源(Reactive Streams 术语中的Publisher,Akka Streams 中实现为Source)保证绝不向任何Subscriber发射超过其已接收总需求数量的元素

注意:Reactive Streams 规范以Publisher/Subscriber定义协议,但这两个类型不是面向用户的 API,而是不同 Reactive Streams 实现之间的底层构件。Akka Streams 将其实现为SourceFlow(对应规范中的Processor)、Sink,不直接暴露 Reactive Streams 接口。需要与其他响应式流库集成时,参见 与 Reactive Streams 集成。

这种背压工作模式可以通俗地称为"动态 push / pull 模式":根据下游能否跟上上游生产速率,在基于 push 与基于 pull 的背压模型之间切换。

4.1 慢 Publisher、快 Subscriber(push 模式)

这是理想情况——无需放慢 Publisher。但信号速率很少恒定,可能随时变化,突然变成"Subscriber 比 Publisher 慢",因此背压协议在这种场景下也必须保持启用,同时不希望为这个安全网付出过高代价。

协议通过 Subscriber 异步向 Publisher 发送Request(n)信号解决:协议保证 Publisher 发射的元素不会超过已声明的 demand。由于当前 Subscriber 更快,它会以更高频率发送 Request 信号(也可能批量合并 demand,一次请求多个元素)。这意味着 Publisher 发射传入元素时几乎不需要等待(被背压)。此场景实际运行在push 模式:Publisher 能多快就多快地产出,因为待处理的 demand 会在发射元素时被"恰好及时"地补充。

4.2 快 Publisher、慢 Subscriber(pull 模式)

此时必须对 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方法)把算子标记为以异步方式与其下游通信。文档示例(Scala):

Source(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.queueSource.actorRefSource.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 ordering)

Akka Streams 中几乎所有计算算子都保持输入元素的顺序:若输入{IA1,...,IAn}"引起"输出{OA1,...,OAk},输入{IB1,...,IBm}"引起"输出{OB1,...,OBl},且所有IAi都先于所有IBi,那么OAi也先于OBi

这一性质对mapAsync这类异步算子也成立;但存在不保序的版本mapAsyncUnordered,它不保持这种顺序。

然而,处理多输入流的连接件(如Merge不保证来自不同输入端口的元素的输出顺序——merge 类操作可能先发射Ai再发射Bi,顺序由内部逻辑决定。而Zip这类专门的算子保证输出顺序,因为每个输出元素都依赖所有上游元素已被发出——因此 zip 场景的顺序由这一性质定义。

如果需要在 fan-in 场景下对元素发射顺序做细粒度控制,可以考虑MergePreferredMergePrioritized,或者自定义GraphStage——它给你对合并方式的完全控制权。

7. Actor Materializer 生命周期

Materializer负责把流蓝图变成运行中的流并产出"物化值"。一个 ActorSystem 级别的Materializer由 AkkaExtensionSystemMaterializer提供——Scala 通过隐式ActorSystem,Java 通过向各种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,将其生命周期绑定到该 Actor(Scala,见 FlowDocSpec.scala):

final class RunWithMyself extends Actor { implicit val mat: Materializer = Materializer(context) Source.maybe.runWith(Sink.onComplete { case Success(done) => println(s"Completed: $done") case Failure(ex) => println(s"Failed: ${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 已完成全部职责、想尽早停止。此时应把系统物化器传入 Actor:

final class RunForever(implicit val mat: Materializer) extends Actor { Source.maybe.runWith(Sink.onComplete { case Success(done) => println(s"Completed: $done") case Failure(ex) => println(s"Failed: ${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:自定义算子与精细控制流行为;
  • 流图与图 DSL:fan-in/fan-out 拓扑、图的物化值访问与环的死锁处理;
  • 与 Reactive Streams 集成:跨实现互操作;
  • 物化器相关实现可阅读 Source.scala 与 SystemMaterializer.scala。
  • 后端
  • 并发编程
  • 异步编程

【免费下载链接】akka-core

A 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
点击查看免费下载
上一篇:Sublime Text编码转换终极指南:告别乱码的完整解决方案
下一篇:小米手表表盘设计终极指南:Mi-Create免费工具完全教程

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询