运行管道
将图物化并执行。
运行管道 是 CoddyKit 上的免费 Scala for Backend Engineering & Functional Programming 课时。 这是第 4 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 Scala for Backend Engineering & Functional Programming 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 Scala for Backend Engineering & Functional Programming 课程共包含 4 节课。
从蓝图到执行
到目前为止,流水线还只是一个纯粹的蓝图。物化是将该蓝图转换为实际移动数据的运行中 actor 的过程。
除非您明确运行该图,否则什么也不会发生;这正是 Akka Streams 能够组合和复用的原因。
ActorSystem
物化需要一个 ActorSystem,它为流的各个阶段提供所依赖的线程和调度器。在现代 Akka 中,该系统还会充当隐式物化器。
一个 ActorSystem 通常可以服务于整个应用程序和许多并发流。
import akka.actor.ActorSystem
implicit val system: ActorSystem =
ActorSystem("data-pipeline")
import system.dispatcher // ExecutionContextrunWith
运行 Source 最直接的方式是使用 runWith:它会附加一个 Sink,并在一步内完成物化,返回该 Sink 的物化值。
这里的结果是一个 Future[Int],流完成后会得到总和。
import akka.stream.scaladsl.{Source, Sink}
import scala.concurrent.Future
val total: Future[Int] =
Source(1 to 100).runWith(Sink.fold(0)(_ + _))在 RunnableGraph 上运行
如果您已经使用 to 或 toMat 构建了一个闭合的 RunnableGraph,请调用 run() 对其进行物化。返回值就是该图所保留的物化值。
这样可以清晰地将流水线构建与执行分离。
import akka.stream.scaladsl.Keep
import scala.concurrent.Future
val graph =
Source(1 to 100)
.toMat(Sink.fold(0)(_ + _))(Keep.right)
val result: Future[Int] = graph.run()便捷的运行运算符
Source 提供了几个快捷方法:runForeach、runFold 和 runReduce 会分别附加相应的 Sink 并立即运行。
对于 Source 上常见的终结操作,它们写法简洁。
import scala.concurrent.Future
val printed: Future[akka.Done] =
Source(1 to 10).runForeach(println)
val sum: Future[Int] =
Source(1 to 10).runFold(0)(_ + _)处理结果 Future
终结 Sink 会返回一个 Future,并在流结束或失败时完成。使用 onComplete 注册回调,以便对成功或错误做出响应。
将 ActorSystem 的调度器用作这些回调的隐式 ExecutionContext。
import scala.util.{Success, Failure}
total.onComplete {
case Success(value) => println(s"Sum = $value")
case Failure(ex) => println(s"Failed: ${ex.getMessage}")
}一条真实的流水线
典型的数据管道从源读取数据,使用流进行转换,通过 mapAsync 执行异步 I/O,使用 grouped 进行批处理,最后写入接收器。
每个阶段都很小,整个管道只需一次 run 即可完成物化。
val done =
lineSource
.map(parse)
.mapAsync(4)(validate)
.grouped(500)
.runWith(bulkWriteSink)重启失败的流
为了提高容错能力,请使用 RestartSource.withBackoff 包装源或流,这样瞬时故障(例如连接断开)就会触发带指数退避的自动重启。
这样无需手动监控,也能让长期运行的数据摄取管道持续工作。
import akka.stream.scaladsl.RestartSource
import akka.stream.RestartSettings
import scala.concurrent.duration._
val resilient = RestartSource.withBackoff(
RestartSettings(1.second, 30.seconds, 0.2))(() => flakySource)使用 KillSwitch 优雅地关闭
KillSwitch 允许外部代码干净地停止正在运行的流。请通过 viaMat 插入 KillSwitches.single,并保留其物化值,以便稍后调用 shutdown()。
对于必须在应用程序关闭时停止的长期运行流,这是必不可少的。
import akka.stream.{KillSwitches, KillSwitch}
import akka.stream.scaladsl.Keep
val (switch, done) =
source
.viaMat(KillSwitches.single)(Keep.right)
.toMat(Sink.ignore)(Keep.both)
.run()
// later: switch.shutdown()释放资源
应用程序退出时,请终止 ActorSystem 以释放其线程。请将 terminate 调用接在流的完成 Future 之后,从而按顺序完成关闭。
泄漏 ActorSystem 会使 JVM 保持运行,并使资源持续处于打开状态。
done.onComplete { _ =>
system.terminate()
}重复使用物化器
多次物化同一个蓝图会创建彼此独立的运行中流,同时共享 ActorSystem 的资源。蓝图本身保持不可变且没有副作用。
因此,您可以安全地定义一次管道,并在每个传入任务到达时按需运行它。
val blueprint =
Source(1 to 5).toMat(Sink.seq)(Keep.right)
val run1 = blueprint.run()
val run2 = blueprint.run() // independent execution快速检查
请思考一个流要实际处理元素,需要满足哪些条件。
回顾
运行管道意味着通过 run、runWith 或便捷运算符,使用一个 ActorSystem 将蓝图物化;每种方式都会返回一个异步结果。
您已经了解了真实的多阶段管道、带退避的自动重启、通过 KillSwitch 实现的优雅关闭、使用 system.terminate() 清理资源,以及如何在彼此独立的运行中安全地重复使用不可变蓝图。
用 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 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 4 节课,共 4 节。
「运行管道」课时需要多长时间?
大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。
我能在这节 Scala for Backend Engineering & Functional Programming 课中编写并运行代码吗?
能。每节 Scala for Backend Engineering & Functional Programming 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。