パイプラインを実行する
グラフをマテリアライズして実行します。
「パイプラインを実行する」はCoddyKit上の無料Scala for Backend Engineering & Functional Programmingレッスンです。 これはレッスン4/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはScala for Backend Engineering & Functional Programming学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 Scala for Backend Engineering & Functional Programmingコースには全4レッスンが含まれています。
ブループリントから実行へ
ここまで、パイプラインは純粋なブループリントでした。マテリアライゼーションとは、そのブループリントを実際にデータを移動させる実行中のアクターへ変換するプロセスです。
グラフを明示的に実行するまで何も起こりません。この仕組みによって、Akka Streamsは組み合わせやすく、再利用しやすくなっています。
ActorSystem
マテリアライゼーションにはActorSystemが必要です。ActorSystemは、ストリームのステージを支えるスレッドとディスパッチャーを提供します。最新のAkkaでは、システム自体が暗黙のマテリアライザーとしても機能します。
通常、1つのActorSystemでアプリケーション全体と多数の並行ストリームを処理します。
import akka.actor.ActorSystem
implicit val system: ActorSystem =
ActorSystem("data-pipeline")
import system.dispatcher // ExecutionContextrunWith
Sourceを実行する最も直接的な方法はrunWithです。Sinkを接続して1ステップでマテリアライズし、その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にコールバックを登録すると、成功またはエラーに応じて処理できます。
これらのコールバックの暗黙のExecutionContextには、ActorSystemのディスパッチャーを使います。
import scala.util.{Success, Failure}
total.onComplete {
case Success(value) => println(s"Sum = $value")
case Failure(ex) => println(s"Failed: ${ex.getMessage}")
}現実的なパイプライン
一般的なデータパイプラインでは、Source から読み取り、Flow で変換し、mapAsync で非同期 I/O を実行し、grouped でバッチ化して、Sink に書き込みます。
各ステージは小さく、全体は単一の run でマテリアライズされます。
val done =
lineSource
.map(parse)
.mapAsync(4)(validate)
.grouped(500)
.runWith(bulkWriteSink)失敗したストリームの再起動
堅牢性を高めるには、Source または Flow を 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 を終了させます。ストリームの完了を表す Future の後に terminate 呼び出しを連結すると、順序正しくシャットダウンできます。
ActorSystem を解放しないと JVM が終了せず、リソースも開いたままになります。
done.onComplete { _ =>
system.terminate()
}Materializer の再利用
同じブループリントを複数回マテリアライズすると、ActorSystem のリソースを共有する、独立して実行されるストリームが作成されます。ブループリント自体は不変で、副作用もありません。
そのため、パイプラインを一度定義しておき、受信するジョブごとにオンデマンドで実行しても安全です。
val blueprint =
Source(1 to 5).toMat(Sink.seq)(Keep.right)
val run1 = blueprint.run()
val run2 = blueprint.run() // independent executionクイックチェック
ストリームが実際に要素を処理するために何が必要かを考えてみましょう。
まとめ
パイプラインを実行するとは、ActorSystem を使ってブループリントを run、runWith、または便利な演算子でマテリアライズすることです。いずれも Future の結果を返します。
現実的な複数ステージのパイプライン、バックオフ付きの自動再起動、KillSwitch によるグレースフルシャットダウン、system.terminate() によるリソースのクリーンアップ、そして不変のブループリントを独立した実行間で安全に再利用する方法を学びました。
AI チューターと学ぶ Scala — 無料
ブラウザでリアルコードを書いて実行し、24/7 の AI チューターから瞬時にサポートを受け、ウェブまたはアプリで続きから学習できます。
- コース
- 39
- レッスン
- 143
よくある質問
「パイプラインを実行する」レッスンは無料ですか?
はい。「パイプラインを実行する」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、Scala for Backend Engineering & Functional Programmingコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 Scala for Backend Engineering & Functional Programmingコースには全4レッスンが含まれています。
「パイプラインを実行する」で何を学びますか?
グラフをマテリアライズして実行します。 ブラウザで直接実行するハンズオンコードでScala for Backend Engineering & Functional Programmingを演習し、24時間対応の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フィードバックを取得できます。ローカル設定は不要です。
このコースのすべてのレッスン
- Source、Flow、Sink
- ストリームを変換する
- バックプレッシャー
- パイプラインを実行する