ストリームを変換する
流れるデータをmapとfilterで処理します。
「ストリームを変換する」はCoddyKit上の無料Scala for Backend Engineering & Functional Programmingレッスンです。 これはレッスン2/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはScala for Backend Engineering & Functional Programming学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 Scala for Backend Engineering & Functional Programmingコースには全4レッスンが含まれています。
変換としてのオペレーター
Akka Streamsには、ScalaのコレクションAPIに似た豊富なSourceおよびFlow用オペレーターが用意されています。これらは非同期で実行され、バックプレッシャーにも対応します。
各オペレーターは新しいブループリントを返すため、ストリームを実行する前に変換を宣言的に組み合わせられます。
mapとfilter
mapは各要素に同期関数を適用し、filterは述語を満たさない要素を除外します。これらは要素単位の変換で最も基本となるオペレーターです。
どちらも順序を維持し、完了と失敗を下流へ伝播します。
val flow =
Flow[Int]
.filter(_ % 2 == 0)
.map(n => n * n)1対多の変換に使うmapConcat
1つの入力から複数の出力を生成したい場合は、mapConcatを使います。これはIterableを返す関数を受け取り、その結果をストリームに平坦化します。
空のコレクションを返すと、その要素は実質的に除外されます。
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)は重なり合うウィンドウを出力します。
バッチ処理により、I/O負荷の高いパイプラインで要素ごとにかかるオーバーヘッドを削減できます。
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はマテリアライズごとに新しい状態を作成し、出力のIterableを返します。
ストリームの実行間で状態を共有せずに、カウンターやバッファーを保持する安全な方法です。
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))Flowの組み合わせ
小さなFlowはviaで組み合わせて、1つの再利用可能な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、1対多のmapConcat、バッチ化のgrouped、scanとfoldによる累積、そしてmapAsyncによる非同期処理です。
状態を持つ変換、throttleなどの時間ベースのオペレーター、エラーに対する監督戦略、viaによるFlowの組み合わせについても学びました。次は、バックプレッシャーによってこれらのステージを安全に保つ方法です。
よくある質問
「ストリームを変換する」レッスンは無料ですか?
はい。「ストリームを変換する」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、Scala for Backend Engineering & Functional Programmingコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 Scala for Backend Engineering & Functional Programmingコースには全4レッスンが含まれています。
「ストリームを変換する」で何を学びますか?
流れるデータをmapとfilterで処理します。 ブラウザで直接実行するハンズオンコードでScala for Backend Engineering & Functional Programmingを演習し、24時間対応の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フィードバックを取得できます。ローカル設定は不要です。
このコースのすべてのレッスン
- Source、Flow、Sink
- ストリームを変換する
- バックプレッシャー
- パイプラインを実行する