Source、Flow、Sink
ストリーミングを構成する基本要素です。
「Source、Flow、Sink」はCoddyKit上の無料Scala for Backend Engineering & Functional Programmingレッスンです。 これはレッスン1/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはScala for Backend Engineering & Functional Programming学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 Scala for Backend Engineering & Functional Programmingコースには全4レッスンが含まれています。
3つの基本構成要素
Akka Streamsでは、データパイプラインを処理ステージのグラフとしてモデル化します。中核となる3つの線形ステージは、要素を生成するSource、要素を変換するFlow、要素を消費するSinkです。
Sourceには出力が1つ、Sinkには入力が1つ、Flowには入力と出力がそれぞれ1つだけあります。これらを接続することで、いつ実行するかではなく、何を実行するかを記述します。
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は1つの結果に集約します。
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によって整数を2倍し、文字列に変換します。
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)は、出力型がFlowの出力型と一致する新しいSourceを生成します。
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]は、型を変換するFlowを間に挟まなければ、Sink[String]に接続できません。
この静的チェックにより、パイプラインの接続ミスを実行前に検出できます。
// 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をマテリアライズします。デフォルトでは、ステージを結合すると左端のマテリアライズ値が保持されます(通常のSourceでは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)再利用可能なコンポーネント
Source、Flow、Sinkは不変値なので、一度定義すれば複数のパイプラインで再利用できます。これにより、小さく名前の付いた処理ステージのライブラリを構築しやすくなります。
解析用に定義したFlowを、ファイル用パイプラインとHTTP用パイプラインの両方に組み込めます。
val parse: Flow[String, Int, akka.NotUsed] =
Flow[String].map(_.trim.toInt)
val fromFile = lines.via(parse)
val fromHttp = requestBody.via(parse)よく使うSourceコンストラクター
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コンストラクター
同様に、Sinkには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理解度チェック
線形ステージ型についての理解を確認しましょう。
まとめ
3つの線形ビルディングブロックを学びました。Sourceはデータを生成し、Flowは変換し、Sinkは消費します。これらは不変で再利用可能なブループリントであり、viaとtoで接続します。
SourceとSinkを接続すると、マテリアライズ値を持つRunnableGraphが生成されます。ただし、実行するまでデータは移動しません。次は、より豊富なオペレーターを使ってストリームを変換します。
よくある質問
「Source、Flow、Sink」レッスンは無料ですか?
はい。「Source、Flow、Sink」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、Scala for Backend Engineering & Functional Programmingコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 Scala for Backend Engineering & Functional Programmingコースには全4レッスンが含まれています。
「Source、Flow、Sink」で何を学びますか?
ストリーミングを構成する基本要素です。 ブラウザで直接実行するハンズオンコードでScala for Backend Engineering & Functional Programmingを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。
Scala for Backend Engineering & Functional Programmingを始めるのに経験は必要ですか?
事前経験は必要ありません。CoddyKitのScala for Backend Engineering & Functional Programmingは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン1/4です。
「Source、Flow、Sink」レッスンにはどのくらい時間がかかりますか?
ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。
このScala for Backend Engineering & Functional Programmingレッスンでコードを書いて実行できますか?
はい。すべてのScala for Backend Engineering & Functional Programmingレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。
このコースのすべてのレッスン
- Source、Flow、Sink
- ストリームを変換する
- バックプレッシャー
- パイプラインを実行する