Kotlin: Kotlin Flow詳解

最終更新:2026-08-26

Flowは、コルーチン向けのストリーミングAPIです。コールドストリームはオンデマンドでデータを配信し、バックプレッシャー機能が組み込まれており、オペレーターは組み合わせ可能です。CharlieはFlowを使用して、リアルタイムの注文処理パイプラインを構築しています。

1. 学習内容


2. 本物の建築家の物語

(1) 課題:リアルタイム注文処理のボトルネック

CharlieのOrderProcessorは、1秒あたり5,000件の注文イベントをリアルタイムで処理する必要があります。List を使ったバッチ処理ではレイテンシが高くなり、コールバックのネストは保守性が悪く、RxJavaは習得に時間がかかります。

(2) フロー・ソリューション

KOTLIN
// Flow: cold stream, on-demand emission, backpressure built-in
fun orderEventStream(): Flow<OrderEvent> = flow {
    while (true) {
        val event = pollOrderEvent()  // Non-blocking poll
        emit(event)                   // Emit to downstream
    }
}

// Compose pipeline: filter -> transform -> collect
orderEventStream()
    .filter { it.total > 100 }
    .map { enrich(it) }
    .collect { dispatch(it) }

Flowは、コルーチンネイティブなストリーミング機能です。オンデマンドでのコールドエミッション、組み込みのバックプレッシャー、組み合わせ可能なオペレーターを備え、RxJavaよりもはるかに習得が容易です。


3. フローの基礎

(1) フローの作成

KOTLIN
// 1. flow builder
fun numberFlow(): Flow<Int> = flow {
    for (i in 1..5) {
        emit(i)  // Emit value to collector
        delay(100)  // Simulate work
    }
}

// 2. flowOf: from fixed values
val orderFlow = flowOf(
    Order("ORD-001", 299.99),
    Order("ORD-002", 1_500.00)
)

// 3. asFlow: from Iterable/Sequence
val orderList = listOf(Order("ORD-001", 299.99))
val orderFlow2 = orderList.asFlow()

// 4. channelFlow: for complex concurrent emission
fun concurrentFlow(): Flow<Int> = channelFlow {
    launch { send(1) }
    launch { send(2) }
}

(2) 低温流の挙動

KOTLIN
// Cold flow: each collector triggers new emission
val coldFlow = flow {
    println("Emitting...")  // This runs EACH time collect is called
    emit(1)
    emit(2)
}

coldFlow.collect { println("Collector 1: $it") }  // Emits 1, 2
coldFlow.collect { println("Collector 2: $it") }  // Emits 1, 2 AGAIN

(3) 低温流動と高温流動

次元 コールドフロー (Flow) ホットフロー (SharedFlow/StateFlow)
トリガー 収集された場合にのみ発火 コレクターとは独立して実行される
データ 各コレクターは全データを受け取る 新しいコレクターはそれ以降のデータのみを確認できる
キャッシュ なし 設定可能な再生
ライフサイクル コレクターに従う 独立
類推 ビデオ・オン・デマンド ライブ配信

4. 中間操作

中間操作は「コールド」です。つまり、エミッションを発生させることはなく、変換ロジックを定義するだけです。

(1) map / filter / transform

KOTLIN
// map: transform each value
val totals = orderFlow.map { it.total }

// filter: keep matching values
val highValue = orderFlow.filter { it.total > 1_000 }

// transform: more flexible than map (can emit multiple values)
val expanded = orderFlow.transform { order ->
    emit("Order: ${order.id}")
    emit("Total: \$${order.total} USD")
}

(2) テイク/ドロップ/デバウンス

KOTLIN
// take: take first N values
val first3 = orderFlow.take(3)

// drop: skip first N values
val after5 = orderFlow.drop(5)

// debounce: emit only after silence period
val debounced = searchFlow.debounce(300)  // Wait 300ms of no new values

(3) 中間操作カテゴリ

カテゴリ 演算子 関数
変換 map / transform 値ごとの変換
フィルタ filter / filterNot / filterIsInstance 条件付きフィルタリング
フラット化 flatMapConcat / flatMapMerge / flatMapLatest ネストされたストリームのフラット化
切り捨て take / drop / takeWhile 最初の N 個を取り出す/スキップ
タイミング debounce / sample / throttleFirst 時間ベース制御
結合 zip / combine マルチストリーム結合

5. ターミナル運営

ターミナルでの作業により、実際の排出および回収が行われます。

(1) collect / toList / first

KOTLIN
// collect: consume each value
orderFlow.collect { order -> println(order.id) }

// toList: collect to list
val allOrders = orderFlow.toList()

// first: get first value
val firstOrder = orderFlow.first()

// single: expect exactly one value
val onlyOrder = orderFlow.single()

(2) ターミナルの運用区分

操作 戻り値の型 説明
collect 単位 各値あたりの消費量
toList / toSet List / Set コレクションに追加
first / firstOrNull T / T? 最初の値を取得
single / singleOrNull T / T? 一意の値を取得
count Int 要素の個数
fold / reduce R / T 累積

6. 背圧への対応

(1) 背圧の問題

KOTLIN
// Producer is faster than consumer
val fastFlow = flow {
    repeat(10) {
        emit(it)           // Fast: 1ms per emit
        delay(1)
    }
}

// Slow consumer causes backpressure
fastFlow.collect {
    delay(100)            // Slow: 100ms per process
    println("Processed: $it")
}
// By default: producer waits for consumer (suspend)

(2) 背圧対策

KOTLIN
// buffer: run producer and consumer in parallel
fastFlow.buffer(10)  // Buffer up to 10 values
    .collect { delay(100); println("Buffered: $it") }

// conflate: skip intermediate values, keep latest
fastFlow.conflate()
    .collect { delay(100); println("Latest: $it") }

// collectLatest: cancel previous collection when new value arrives
fastFlow.collectLatest { value ->
    delay(100)
    println("Latest processed: $value")
}

(3) 背圧戦略の比較

戦略 行動 データ損失 ユースケース
デフォルト プロデューサーがコンシューマーを待機 データ損失なし 完全性の保証
buffer N個の値をバッファに格納 バッファがいっぱいになると一時停止 短時間の生産量の急増
conflate 中間値をスキップし、最新値のみ保持 中間値が失われる リアルタイム表示
collectLatest 以前の処理をキャンセルし、最新の処理を実行 以前の処理が中断されました 検索候補

7. フロー・コールド・ストリーム・パイプライン

100%
flowchart LR
    A[Source<br/>flow emit] --> B[filter<br/>total > 100]
    B --> C[map<br/>enrich]
    C --> D[debounce<br/>300ms]
    D --> E[collect<br/>dispatch]

8. 完全な例:OrderProcessor イベントストリーム

▶ サンプル:注文イベントストリーム

KOTLIN
// ============================================
// OrderProcessor - Event Stream with Flow
// Feature: Real-time order event processing
// ============================================

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*

data class Order(val id: String, val total: Double, val status: String, val customer: String)

sealed class OrderEvent {
    data class Created(val order: Order) : OrderEvent()
    data class Updated(val orderId: String, val newStatus: String) : OrderEvent()
    data class Cancelled(val orderId: String, val reason: String) : OrderEvent()
}

// Simulate event source
fun orderEventStream(): Flow<OrderEvent> = flow {
    val events = listOf(
        OrderEvent.Created(Order("ORD-001", 299.99, "PENDING", "Alice")),
        OrderEvent.Created(Order("ORD-002", 15_000.00, "PENDING", "Bob")),
        OrderEvent.Updated("ORD-001", "CONFIRMED"),
        OrderEvent.Created(Order("ORD-003", 2_500.00, "PENDING", "Charlie")),
        OrderEvent.Cancelled("ORD-003", "Out of stock"),
        OrderEvent.Updated("ORD-002", "CONFIRMED"),
        OrderEvent.Created(Order("ORD-004", 8_900.00, "PENDING", "Bob")),
        OrderEvent.Updated("ORD-002", "SHIPPED")
    )
    events.forEach { event ->
        emit(event)
        delay(100)
    }
}

// Enrich created orders
fun Order.enrich() = copy(status = "ENRICHED")

fun main() = runBlocking {
    println("=== Order Event Stream ===")

    // Pipeline: filter created -> enrich -> collect
    orderEventStream()
        .onStart { println("Stream started") }
        .onEach { println("  Received: $it") }
        .filter { it is OrderEvent.Created }
        .map { (it as OrderEvent.Created).order }
        .filter { it.total > 1_000 }
        .map { it.enrich() }
        .onEach { println("  >>> Enriched high-value: ${it.id} (\$${it.total} USD)") }
        .onCompletion { println("Stream completed") }
        .collect()

    // Separate pipeline: status updates
    println("\n=== Status Updates ===")
    orderEventStream()
        .filter { it is OrderEvent.Updated }
        .map { it as OrderEvent.Updated }
        .collect { event ->
            println("  ${event.orderId} -> ${event.newStatus}")
        }

    // Aggregate: count events by type
    println("\n=== Event Counts ===")
    val counts = orderEventStream()
        .groupBy { it::class.simpleName ?: "Unknown" }
        .map { (type, flow) ->
            type to flow.count()
        }
        .toList()
        .associate { it }

    counts.forEach { (type, count) ->
        println("  $type: $count events")
    }
}

出力:

TEXT 📖 参照専用
=== Order Event Stream ===
Stream started
  Received: Created(order=Order(ORD-001, 299.99, PENDING, Alice))
  Received: Created(order=Order(ORD-002, 15000.0, PENDING, Bob))
  >>> Enriched high-value: ORD-002 ($15000.0 USD)
  Received: Updated(orderId=ORD-001, newStatus=CONFIRMED)
  Received: Created(order=Order(ORD-003, 2500.0, PENDING, Charlie))
  >>> Enriched high-value: ORD-003 ($2500.0 USD)
  Received: Cancelled(orderId=ORD-003, reason=Out of stock)
  Received: Updated(orderId=ORD-002, newStatus=CONFIRMED)
  Received: Created(order=Order(ORD-004, 8900.0, PENDING, Bob))
  >>> Enriched high-value: ORD-004 ($8900.0 USD)
  Received: Updated(orderId=ORD-002, newStatus=SHIPPED)
Stream completed

=== Status Updates ===
  ORD-001 -> CONFIRMED
  ORD-002 -> CONFIRMED
  ORD-002 -> SHIPPED

=== Event Counts ===
  Created: 4 events
  Updated: 3 events
  Cancelled: 1 events

❓ よくある質問

Q FlowとSequenceの違いは何ですか?
A Flowは非同期(サスペンドに対応)で、Sequenceは同期です。Flowは非同期のデータソース(ネットワーク/データベース)に適しており、Sequenceは同期型のビッグデータ処理に適しています。
Q FlowとRxJavaのObservableの違いは何ですか?
A Flowはコルーチンネイティブで、コールドストリームファーストを重視しており、習得のハードルが低いのが特徴です。RxJavaはより成熟していますが、複雑さも増します。新規プロジェクトではFlowを選択し、既存のプロジェクトではRxJavaと併用することをお勧めします。
Q コールドフローはいつホットフローになるのですか?
A shareIn を使用するとコールドフローを SharedFlow に変換でき、stateIn を使用すると StateFlow に変換できます。これにより、複数のコレクターが同じデータソースを共有できるようになります。
Q flatMapConcat、flatMapMerge、flatMapLatestのどれを選べばいいですか?
A flatMapConcat は順次処理(デフォルト)、flatMapMerge は並列処理(スループットが高い)、flatMapLatest は古いストリームを破棄します(検索などのリアルタイムシナリオ向け)。
Q Flowはデフォルトでバックプレッシャーをどのように処理しますか?
A デフォルトの挙動は「プロデューサーがコンシューマーを待つ」というものです。— emit はsuspend関数であり、コンシューマーの処理が遅い場合、プロデューサーは自動的にサスペンドされます。これが最も安全なデフォルト設定です。
Q collect はsuspend関数ですか?
A はい。collect はsuspend関数であり、コルーチンまたはsuspend関数からのみ呼び出すことができます。Flow が完了するまで実行を中断します。

📖 まとめ


📝 練習問題

  1. 初心者 (⭐): flow { } を使用して 1 から 10 までの数値を出力するフローを作成し、その中から偶数をフィルタリングして表示してください。ヒント: flow { for (i in 1..10) emit(i) }.filter { }
  2. 中級 (⭐⭐): Flow を使用して注文イベントストリームをシミュレートし、filter + map + collect を使用して、金額が 1000 を超える注文を処理してください。ヒント:第 8 節を参照してください。
  3. 上級 (⭐⭐⭐): debounce を使用して検索ストリームを実装してください。ユーザー入力をシミュレートし、300ms間新しい入力がない場合にのみ検索をトリガーするようにしてください。ヒント: debounce(300)

← 前へ | 次へ →

Web-Tutorial.com

Web-Tutorial 技術チーム

複数の開発者によって共同維持されているプログラミングチュートリアルプラットフォーム。各チュートリアルは専門分野の開発者が執筆・レビューしています。正確で信頼性の高いコンテンツを目指しています — 問題を見つけた場合はお知らせください。

100%