flowOnとbuffer:コンテキストとバックプレッシャー
flowOnで発行コンテキストを変更し、bufferでバックプレッシャーに対応します。
「flowOnとbuffer:コンテキストとバックプレッシャー」はCoddyKit上の無料Kotlin Academyレッスンです。 これはレッスン4/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはKotlin Academy学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 Kotlin Academyコースには全4レッスンが含まれています。
Flowのコンテキスト
デフォルトでは、フローはcollectを呼び出したコルーチンのコンテキストで実行されます。変更しない限り、プロデューサーとコンシューマーのディスパッチャーは同じです。
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
flow {
println("emit on: ${Thread.currentThread().name}")
emit(1)
}.collect {
println("collect on: ${Thread.currentThread().name}")
}
}flowOnで上流のコンテキストを変更する
flowOn(dispatcher)は、上流のフロー(チェーン内でその位置より上にあるすべての処理)を指定したディスパッチャーで実行します。一方、収集処理は呼び出し元のディスパッチャーで実行されます。
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
flow {
println("emit: ${Thread.currentThread().name}")
emit(1)
}.map {
println("map: ${Thread.currentThread().name}")
it * 2
}.flowOn(Dispatchers.Default) // above runs on Default
.collect {
println("collect: ${Thread.currentThread().name}")
}
}チェーン内で複数のflowOnを使用する
flowOnは複数回使用できます。それぞれのflowOnは、直前のflowOnまでの、その直上にあるオペレーターに影響します。
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
flow { emit(1) }
.map { it + 1 }.flowOn(Dispatchers.IO) // map runs on IO
.map { it * 2 }.flowOn(Dispatchers.Default) // this map runs on Default
.collect { println("Result: $it") } // collect on Main (runBlocking)
}バックプレッシャーの問題
プロデューサーがコンシューマーの処理速度より速く値を発行すると、値がキューに蓄積されます。バッファリングしない場合、プロデューサーは待機することになります。
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
val time = System.currentTimeMillis()
flow {
repeat(3) { i ->
delay(100) // fast producer
emit(i)
}
}.collect {
delay(300) // slow consumer
println("Got $it in ${System.currentTimeMillis() - time}ms")
}
}buffer()オペレーター
buffer()は、プロデューサーとコンシューマーを別々のコルーチンで並行して実行し、発行された値をチャネルにバッファリングします。プロデューサーはコンシューマーを待機しません。
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
val time = System.currentTimeMillis()
flow {
repeat(3) { i -> delay(100); emit(i) }
}.buffer() // producer and consumer run concurrently
.collect {
delay(300)
println("Got $it in ${System.currentTimeMillis() - time}ms")
}
}bufferの容量
buffer(capacity)はチャネルのバッファサイズを設定します。バッファが満杯になると、プロデューサーは一時停止します(バックプレッシャー)。デフォルトの容量は64です。
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlinx.coroutines.channels.Channel
fun main() = runBlocking {
flow { repeat(5) { emit(it) } }
.buffer(Channel.RENDEZVOUS) // 0: producer waits
// .buffer(Channel.BUFFERED) // default: 64
// .buffer(Channel.UNLIMITED) // unbounded
.collect { delay(50); println(it) }
}最新値だけにするconflate()
conflate()は、コレクターの処理が遅い場合に途中の値を破棄し、最も新しい値の発行だけを保持します。UI状態に便利です。
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
flow {
repeat(5) { i -> emit(i); delay(50) }
}.conflate()
.collect { i ->
delay(150)
println("Collected: $i") // skips some values
}
}遅いコレクター向けのcollectLatest
collectLatestは新しい値が到着すると現在の収集ブロックをキャンセルし、最新値で再開します。
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
flow {
emit(1); delay(50)
emit(2); delay(50)
emit(3)
}.collectLatest { value ->
println("Processing $value")
delay(100) // gets cancelled if new value arrives
println("Done $value") // only prints for last value
}
}flowOn + bufferパターン
flowOnとbufferを組み合わせ、プロデューサーをIO(ネットワーク/ディスク)で実行し、結果をバッファリングして、Mainで収集します。
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun fetchItems(): Flow<String> = flow {
repeat(3) { i ->
delay(100) // simulate IO
emit("item-$i")
}
}.flowOn(Dispatchers.IO).buffer(10)
fun main() = runBlocking {
fetchItems().collect { println("UI: $it") }
}並行プロデューサー向けのchannelFlow
channelFlowはチャネルを基盤とするフローを作成し、ビルダー内から複数のコルーチンが並行して値を発行できるようにします。
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun concurrentFlow(): Flow<Int> = channelFlow {
launch { send(1) }
launch { send(2) }
launch { send(3) }
}
fun main() = runBlocking {
concurrentFlow().collect { println(it) }
}適切な戦略の選択
まとめると、コンテキストの切り替えにはflowOn、スループットの向上にはbuffer、UIで最新値だけを扱うにはconflate、古くなった処理のキャンセルにはcollectLatestを使用します。
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
// Guidelines:
// CPU-heavy production -> flowOn(Dispatchers.Default)
// IO-heavy production -> flowOn(Dispatchers.IO)
// Slow consumer -> buffer()
// UI state updates -> conflate() or StateFlow
// Search/autocomplete -> collectLatest or flatMapLatest
fun main() = runBlocking { println("Choose the right strategy!") }クイックチェック
フローチェーン内でflowOnが影響するのは何ですか。
復習
flowOnは上流の処理を別のディスパッチャーに移します。bufferはプロデューサーとコンシューマーを分離してスループットを向上させます。conflateは途中の値を破棄します。collectLatestは新しい値が到着すると遅い処理をキャンセルします。
よくある質問
「flowOnとbuffer:コンテキストとバックプレッシャー」レッスンは無料ですか?
はい。「flowOnとbuffer:コンテキストとバックプレッシャー」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、Kotlin Academyコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 Kotlin Academyコースには全4レッスンが含まれています。
「flowOnとbuffer:コンテキストとバックプレッシャー」で何を学びますか?
flowOnで発行コンテキストを変更し、bufferでバックプレッシャーに対応します。 ブラウザで直接実行するハンズオンコードでKotlin Academyを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。
Kotlin Academyを始めるのに経験は必要ですか?
事前経験は必要ありません。CoddyKitのKotlin Academyは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン4/4です。
「flowOnとbuffer:コンテキストとバックプレッシャー」レッスンにはどのくらい時間がかかりますか?
ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。
このKotlin Academyレッスンでコードを書いて実行できますか?
はい。すべてのKotlin Academyレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。
このコースのすべてのレッスン
- Flow演算子:map、filter、transform、take
- catchとonCompletion:Flowのエラー処理
- combineとzip:複数のFlowを統合する
- flowOnとbuffer:コンテキストとバックプレッシャー