0Pricing
Kotlin Academy · レッスン

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

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

  1. Flow演算子:map、filter、transform、take
  2. catchとonCompletion:Flowのエラー処理
  3. combineとzip:複数のFlowを統合する
  4. flowOnとbuffer:コンテキストとバックプレッシャー
← Kotlin Academyに戻る