源、流与汇
流处理的构建模块。
源、流与汇 是 CoddyKit 上的免费 Scala for Backend Engineering & Functional Programming 课时。 这是第 1 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 Scala for Backend Engineering & Functional Programming 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 Scala for Backend Engineering & Functional Programming 课程共包含 4 节课。
三个构建模块
Akka Streams 将数据管道建模为由处理阶段组成的图。三个核心线性阶段是Source(生成元素)、Flow(转换元素)和Sink(消费元素)。
Source 有一个输出,Sink 有一个输入,而 Flow 恰好有一个输入和一个输出。将它们连接起来描述的是应该发生什么,而不是何时发生。
定义 Source
Source[Out, Mat]会发出Out类型的元素,并提供Mat类型的物化值。最简单的 Source 来自内存集合或范围。
在流运行之前,Source 只是一个不可变的蓝图,可以自由重用。
import akka.stream.scaladsl.Source
val numbers: Source[Int, akka.NotUsed] =
Source(1 to 100)
val single: Source[String, akka.NotUsed] =
Source.single("hello")定义 Sink
Sink[In, Mat]会消费In类型的元素。物化值通常会保存消费结果,例如在流结束时完成的Future。
Sink.foreach会对每个元素执行副作用;Sink.fold会累积一个结果。
import akka.stream.scaladsl.Sink
import scala.concurrent.Future
val printSink: Sink[Int, Future[akka.Done]] =
Sink.foreach(println)
val sumSink: Sink[Int, Future[Int]] =
Sink.fold(0)(_ + _)定义 Flow
Flow[In, Out, Mat]位于 Source 和 Sink 之间,用于转换每个元素。Flow 本身可以重用,也可以在连接到任何端点之前进行组合。
这里的 Flow 会将整数加倍,并将其转换为字符串。
import akka.stream.scaladsl.Flow
val doubleToString: Flow[Int, String, akka.NotUsed] =
Flow[Int]
.map(_ * 2)
.map(n => s"value=$n")将 Source 连接到 Sink
via运算符会将 Flow 附加到 Source,而to会附加 Sink。使用to将 Source 直接连接到 Sink,会生成一个RunnableGraph:一个封闭且可运行的蓝图。
此时还不会移动任何元素;这里只是在描述拓扑结构。
import akka.stream.scaladsl.{Source, Sink, RunnableGraph}
val graph: RunnableGraph[akka.NotUsed] =
Source(1 to 10).to(Sink.foreach(println))via:插入 Flow
使用 via 将一个 Flow 拼接到流水线中。Source.via(flow) 会生成一个新的 Source,其输出类型与 Flow 的输出类型一致。
链式调用 via 可以将小型、易于测试的 Flow 组件组合成长的转换流水线。
val pipeline =
Source(1 to 10)
.via(Flow[Int].filter(_ % 2 == 0))
.via(Flow[Int].map(_ * 10))
.to(Sink.foreach(println))各阶段的类型安全
编译器会确保每个阶段的输出类型与下一个阶段的输入类型一致。Source[Int] 不能直接连接到 Sink[String],除非中间有一个 Flow 负责转换类型。
这种静态检查能在运行前发现流水线连接错误。
// Source[Int] -> Flow[Int, String] -> Sink[String]
val ok =
Source(1 to 3)
.via(Flow[Int].map(_.toString))
.to(Sink.foreach[String](println))物化值
每个蓝图都携带一个物化值:流运行时生成的句柄。Sink.fold 会物化出一个包含结果的 Future。默认情况下,组合各阶段会保留最左侧的物化值(对于普通源,这是 NotUsed)。
使用 toMat 和 Keep 来选择所需的某一侧的值。
import akka.stream.scaladsl.Keep
import scala.concurrent.Future
val g: RunnableGraph[Future[Int]] =
Source(1 to 100)
.toMat(Sink.fold(0)(_ + _))(Keep.right)可复用组件
由于源、流和汇都是不可变值,您可以定义一次,然后在多个流水线中复用。这有助于构建由小型、命名明确的处理阶段组成的库。
为解析而定义的 Flow 可以同时插入文件流水线和 HTTP 流水线。
val parse: Flow[String, Int, akka.NotUsed] =
Flow[String].map(_.trim.toInt)
val fromFile = lines.via(parse)
val fromHttp = requestBody.via(parse)常用源构造器
Akka Streams 提供了许多 Source 工厂:Source.single、Source.repeat、用于按时间发出元素的 Source.tick、从 Future 创建源的 Source.future,以及 Source.empty。
选择合适的构造器可以明确表达数据生产者的意图。
import scala.concurrent.duration._
val ticks = Source.tick(0.seconds, 1.second, "tick")
val onceF = Source.future(scala.concurrent.Future.successful(42))
val forever = Source.repeat("x")常用汇构造器
同样,汇包括 Sink.head(将第一个元素作为 Future 返回)、Sink.seq(将所有元素收集到 Seq 中)、Sink.ignore(排空并丢弃元素)以及 Sink.last。
对于需要将结果返回给代码的流水线,最常用的是 Sink.seq 和 Sink.fold。
import scala.concurrent.Future
val collect: Sink[Int, Future[Seq[Int]]] = Sink.seq
val firstOne: Sink[Int, Future[Int]] = Sink.head
val drain: Sink[Int, Future[akka.Done]] = Sink.ignore快速检查
测试您对线性阶段类型的理解。
回顾
您学习了三种线性构建块:Source 负责产生数据,Flow 负责转换数据,Sink 负责消费数据。它们是不可变、可复用的蓝图,可以通过 via 和 to 连接。
将 Source 连接到 Sink 会得到一个 RunnableGraph,它携带一个物化值,但在调用 run 之前不会移动任何数据。接下来,您将使用更丰富的运算符转换流。
用 AI 导师学习 Scala — 免费
在浏览器中编写并运行真实代码,获得全天候 AI 导师的即时帮助,并在网页或应用中继续学习。
- 课程
- 39
- 课程
- 143
常见问题解答
「源、流与汇」课时是免费的吗?
是的 — 「源、流与汇」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 Scala for Backend Engineering & Functional Programming 课程的其余内容,请升级到 CoddyKit PRO。 Scala for Backend Engineering & Functional Programming 课程共包含 4 节课。
「源、流与汇」这节课中我会学到什么?
流处理的构建模块。 你通过在浏览器中直接运行的动手代码来练习 Scala for Backend Engineering & Functional Programming,全天候 AI 导师会在你学习这节课的过程中回答你的问题。
学习 Scala for Backend Engineering & Functional Programming 需要有经验吗?
无需任何先前经验。CoddyKit 上的 Scala for Backend Engineering & Functional Programming 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 1 节课,共 4 节。
「源、流与汇」课时需要多长时间?
大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。
我能在这节 Scala for Backend Engineering & Functional Programming 课中编写并运行代码吗?
能。每节 Scala for Backend Engineering & Functional Programming 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。