Scala for Backend Engineering & Functional Programming · 课时

源、流与汇

流处理的构建模块。

第 1 / 4 课13 个步骤

源、流与汇 是 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 反馈 — 无需本地设置。

此课程中的所有课时

  1. 源、流与汇
  2. 转换流
  3. 背压
  4. 运行管道
← 返回 Scala for Backend Engineering & Functional Programming