0Pricing
Scala for Backend Engineering & Functional Programming · レッスン

ストリームを変換する

流れるデータを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フィードバックを取得できます。ローカル設定は不要です。

このコースのすべてのレッスン

  1. Source、Flow、Sink
  2. ストリームを変換する
  3. バックプレッシャー
  4. パイプラインを実行する
← Scala for Backend Engineering & Functional Programmingに戻る