0Pricing
Scala for Backend Engineering & Functional Programming · 课时

转换流

映射并过滤流动的数据。

转换流 是 CoddyKit 上的免费 Scala for Backend Engineering & Functional Programming 课时。 这是第 2 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 Scala for Backend Engineering & Functional Programming 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 Scala for Backend Engineering & Functional Programming 课程共包含 4 节课。

作为转换的运算符

Akka Streams 为 Source 和 Flow 提供了丰富的运算符集。这些运算符类似于 Scala 的集合 API,但以异步方式运行,并遵守背压。

每个运算符都会返回一个新的蓝图,因此可以在流真正运行之前,以声明式方式组合转换。

map 和 filter

map 会对每个元素应用同步函数;filter 会丢弃不满足谓词的元素。它们是逐元素转换中最常用的工具。

两者都会保持顺序,并将完成和失败信号向下游传播。

val flow =
  Flow[Int]
    .filter(_ % 2 == 0)
    .map(n => n * n)

使用 mapConcat 实现一对多转换

当一个输入应产生多个输出时,请使用 mapConcat。它接收一个返回可迭代对象的函数,并将结果展开到流中。

返回空集合即可有效地丢弃该元素。

val explode: Flow[String, String, akka.NotUsed] =
  Flow[String].mapConcat(line => line.split(",").toList)

val words = Source(List("a,b", "c,d,e"))
  .via(explode)

grouped 和 sliding

grouped(n) 会将连续元素分批组成一个最多包含 n 个元素的 Seq,适合批量写入数据库。sliding(n) 会发出相互重叠的窗口。

批处理可以降低以输入输出为主的流水线中的逐元素开销。

val batches: Source[Seq[Int], akka.NotUsed] =
  Source(1 to 1000).grouped(100)

val windows =
  Source(1 to 10).sliding(3, step = 1)

scan 和 fold

scan 会在每个元素之后发出当前累加器的值,从而形成不断变化的状态流。fold 只会在上游完成后发出一次最终累加值。

实时计数器应使用 scan,终结聚合应使用 fold。

val running =
  Source(1 to 5).scan(0)(_ + _) // 0,1,3,6,10,15

val total =
  Source(1 to 5).fold(0)(_ + _) // 15

使用 mapAsync 处理异步工作

mapAsync(parallelism) 会调用一个返回 Future 的函数,并按顺序发出结果,同时最多并发运行 parallelism 个 Future。

当顺序很重要时,可以使用它执行数据库查询或 HTTP 请求等异步调用。

import scala.concurrent.Future

val enriched =
  Flow[UserId]
    .mapAsync(parallelism = 4)(id => lookup(id))

def lookup(id: UserId): Future[User] = ???

mapAsyncUnordered

mapAsyncUnordered 的行为类似于 mapAsync,但会在每个结果完成后立即发出,而不考虑输入顺序。

当下游不关心顺序时,它可以提高吞吐量,因为较慢的 Future 不会再阻塞较快的 Future。

val fast =
  Flow[UserId]
    .mapAsyncUnordered(parallelism = 8)(id => lookup(id))

使用 statefulMapConcat 进行有状态转换

对于需要可变局部状态的逐元素转换,statefulMapConcat 会在每次物化时创建新的状态,并返回一个输出可迭代对象。

这是保存计数器或缓冲区的安全方式,不会在多次流运行之间共享状态。

val withIndex: Flow[String, (Int, String), akka.NotUsed] =
  Flow[String].statefulMapConcat { () =>
    var i = 0
    elem => { i += 1; List((i, elem)) }
  }

基于时间的运算符

流既可以根据内容进行转换,也可以根据时间进行转换。throttle 会限制发出速率,groupedWithin 会按大小或经过的时间进行批处理,takeWithin 会限制持续时间。

这些运算符对于限制外部 API 的访问速率至关重要。

import scala.concurrent.duration._

val limited =
  Source(1 to 1000)
    .throttle(10, 1.second)
    .groupedWithin(100, 500.millis)

处理转换中的错误

默认情况下,运算符内部抛出的异常会导致整个流失败。您也可以使用监督策略,让阶段 resume(丢弃错误元素)或 restart。

在 Flow 上使用 withAttributes 附加该策略。

import akka.stream.{ActorAttributes, Supervision}

val safe =
  Flow[String].map(_.toInt)
    .withAttributes(
      ActorAttributes.supervisionStrategy(_ => Supervision.Resume))

组合 Flows

使用 via 可以将小型 Flow 组合成更大的 Flow,并生成一个可复用的整体 Flow。这样可以让每个转换各司其职,并能够独立测试。

组合后的 Flow 的输入类型取自第一个 Flow,输出类型取自最后一个 Flow。

val parse  = Flow[String].map(_.toInt)
val square = Flow[Int].map(n => n * n)

val parseAndSquare: Flow[String, Int, akka.NotUsed] =
  parse.via(square)

快速检查

思考异步转换及其顺序保证。

回顾

您了解了多种转换运算符:逐元素处理使用 map/filter,一对多处理使用 mapConcat,批处理使用 grouped,累加使用 scan 和 fold,异步工作使用 mapAsync。

您还了解了有状态转换、throttle 等基于时间的运算符、用于处理错误的监督策略,以及如何使用 via 组合 Flow。下一步:了解背压如何保证这些阶段的安全。

常见问题解答

「转换流」课时是免费的吗?

是的 — 「转换流」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 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 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 2 节课,共 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