Kotlin: Kotlin Flow詳解
最終更新:2026-08-26
Flowは、コルーチン向けのストリーミングAPIです。コールドストリームはオンデマンドでデータを配信し、バックプレッシャー機能が組み込まれており、オペレーターは組み合わせ可能です。CharlieはFlowを使用して、リアルタイムの注文処理パイプラインを構築しています。
1. 学習内容
- コールドストリーム対ホットストリーム:
Flowはオンデマンドで配信する - 中間操作:
map/filter/transform/take/debounce - ターミナル操作:
collect/toList/first/single - 背圧対応:
buffer/conflate/collectLatest - Charlieの実践:オーダーイベントのストリーミングパイプライン
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. フロー・コールド・ストリーム・パイプライン
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 が完了するまで実行を中断します。📖 まとめ
- Flowは「コールドストリーム」です:各コレクターは独立してエミッションをトリガーし、データの共有は行われません
- 中間操作(map/filter/transform)はコールド操作である — エミッションをトリガーしない
- ターミナル操作(collect/toList/first)によって、実際のデータフローがトリガーされる
- デフォルトのバックプレッシャー:プロデューサーはコンシューマーを待機します。バッファ/conflate/collectLatest を使用して調整してください。
- フロー演算子の構文はコレクション演算子と一貫しており、習得が容易です
onStart/onEach/onCompletionは、ストリームのライフサイクルフックです
📝 練習問題
- 初心者 (⭐):
flow { }を使用して 1 から 10 までの数値を出力するフローを作成し、その中から偶数をフィルタリングして表示してください。ヒント:flow { for (i in 1..10) emit(i) }.filter { } - 中級 (⭐⭐): Flow を使用して注文イベントストリームをシミュレートし、
filter+map+collectを使用して、金額が 1000 を超える注文を処理してください。ヒント:第 8 節を参照してください。 - 上級 (⭐⭐⭐):
debounceを使用して検索ストリームを実装してください。ユーザー入力をシミュレートし、300ms間新しい入力がない場合にのみ検索をトリガーするようにしてください。ヒント:debounce(300)