Kotlin Coroutines Flows

affaan-m/ECC/pi/core/skills/kotlin-coroutines-flows

作者 affaan-mef648e01899ba3e8dc6371642deaaf64b4477775無授權條款275K 個星標收錄於 2026年10月9日更新於 2026年10月9日儲存庫4 天前更新

Kotlin Coroutines and Flow patterns for Android and KMP — structured concurrency, Flow operators, StateFlow, error handling, and testing. Use when writing coroutines or Flow code on Android or KMP, or debugging cancellation and concurrency.

AI 產生的概覽

針對 Android 與 KMP 的 Kotlin 協程與 Flow 參考模式,涵蓋並行、運算子、取消與測試。

功能
提供結構化並行、Flow 與 StateFlow 用法、調度器、取消以及協程測試的指引與程式碼模式,適用於 Android 和 Kotlin Multiplatform 專案。內容涵蓋平行分解、supervisorScope、debounce 與 retryWhen 等 Flow 運算子、SharedFlow 事件,以及使用 Turbine 與測試調度器進行測試。同時列出應避免的反模式。產出的是說明性指引與範例程式碼,而非檔案或指令碼。
適用情境
適用於撰寫或審查使用協程與 Flow 的 Kotlin 非同步程式碼,或在 Android、KMP 上排查取消、並行與反應式串流行為時。也適合用於測試協程與 Flow。
執行需求
不隨附指令碼或工具,僅為說明性內容。使用範例需具備包含協程的 Kotlin 專案,測試範例涉及 Turbine 與測試調度器。

Kotlin Coroutines & Flows

Patterns for structured concurrency, Flow-based reactive streams, and coroutine testing in Android and Kotlin Multiplatform projects.

When to Activate

  • Writing async code with Kotlin coroutines
  • Using Flow, StateFlow, or SharedFlow for reactive data
  • Handling concurrent operations (parallel loading, debounce, retry)
  • Testing coroutines and Flows
  • Managing coroutine scopes and cancellation

Structured Concurrency

Scope Hierarchy

Application  └── viewModelScope (ViewModel)        └── coroutineScope { } (structured child)              ├── async { } (concurrent task)              └── async { } (concurrent task)

Always use structured concurrency — never GlobalScope:

kotlin
// BADGlobalScope.launch { fetchData() }
// GOOD — scoped to ViewModel lifecycleviewModelScope.launch { fetchData() }
// GOOD — scoped to composable lifecycleLaunchedEffect(key) { fetchData() }

Parallel Decomposition

Use coroutineScope + async for parallel work:

kotlin
suspend fun loadDashboard(): Dashboard = coroutineScope {    val items = async { itemRepository.getRecent() }    val stats = async { statsRepository.getToday() }    val profile = async { userRepository.getCurrent() }    Dashboard(        items = items.await(),        stats = stats.await(),        profile = profile.await()    )}

SupervisorScope

Use supervisorScope when child failures should not cancel siblings:

kotlin
suspend fun syncAll() = supervisorScope {    launch { syncItems() }       // failure here won't cancel syncStats    launch { syncStats() }    launch { syncSettings() }}

Flow Patterns

Cold Flow — One-Shot to Stream Conversion

kotlin
fun observeItems(): Flow<List<Item>> = flow {    // Re-emits whenever the database changes    itemDao.observeAll()        .map { entities -> entities.map { it.toDomain() } }        .collect { emit(it) }}

StateFlow for UI State

kotlin
class DashboardViewModel(    observeProgress: ObserveUserProgressUseCase) : ViewModel() {    val progress: StateFlow<UserProgress> = observeProgress()        .stateIn(            scope = viewModelScope,            started = SharingStarted.WhileSubscribed(5_000),            initialValue = UserProgress.EMPTY        )}

WhileSubscribed(5_000) keeps the upstream active for 5 seconds after the last subscriber leaves — survives configuration changes without restarting.

Combining Multiple Flows

kotlin
val uiState: StateFlow<HomeState> = combine(    itemRepository.observeItems(),    settingsRepository.observeTheme(),    userRepository.observeProfile()) { items, theme, profile ->    HomeState(items = items, theme = theme, profile = profile)}.stateIn(viewModelScope, SharingStarted.WhileSubscribed(5_000), HomeState())

Flow Operators

kotlin
// Debounce search inputsearchQuery    .debounce(300)    .distinctUntilChanged()    .flatMapLatest { query -> repository.search(query) }    .catch { emit(emptyList()) }    .collect { results -> _state.update { it.copy(results = results) } }
// Retry with exponential backofffun fetchWithRetry(): Flow<Data> = flow { emit(api.fetch()) }    .retryWhen { cause, attempt ->        if (cause is IOException && attempt < 3) {            delay(1000L * (1 shl attempt.toInt()))            true        } else {            false        }    }

SharedFlow for One-Time Events

kotlin
class ItemListViewModel : ViewModel() {    private val _effects = MutableSharedFlow<Effect>()    val effects: SharedFlow<Effect> = _effects.asSharedFlow()
    sealed interface Effect {        data class ShowSnackbar(val message: String) : Effect        data class NavigateTo(val route: String) : Effect    }
    private fun deleteItem(id: String) {        viewModelScope.launch {            repository.delete(id)            _effects.emit(Effect.ShowSnackbar("Item deleted"))        }    }}
// Collect in ComposableLaunchedEffect(Unit) {    viewModel.effects.collect { effect ->        when (effect) {            is Effect.ShowSnackbar -> snackbarHostState.showSnackbar(effect.message)            is Effect.NavigateTo -> navController.navigate(effect.route)        }    }}

Dispatchers

kotlin
// CPU-intensive workwithContext(Dispatchers.Default) { parseJson(largePayload) }
// IO-bound workwithContext(Dispatchers.IO) { database.query() }
// Main thread (UI) — default in viewModelScopewithContext(Dispatchers.Main) { updateUi() }

In KMP, use Dispatchers.Default and Dispatchers.Main (available on all platforms). Dispatchers.IO is JVM/Android only — use Dispatchers.Default on other platforms or provide via DI.

Cancellation

Cooperative Cancellation

Long-running loops must check for cancellation:

kotlin
suspend fun processItems(items: List<Item>) = coroutineScope {    for (item in items) {        ensureActive()  // throws CancellationException if cancelled        process(item)    }}

Cleanup with try/finally

kotlin
viewModelScope.launch {    try {        _state.update { it.copy(isLoading = true) }        val data = repository.fetch()        _state.update { it.copy(data = data) }    } finally {        _state.update { it.copy(isLoading = false) }  // always runs, even on cancellation    }}

Testing

Testing StateFlow with Turbine

kotlin
@Testfun `search updates item list`() = runTest {    val fakeRepository = FakeItemRepository().apply { emit(testItems) }    val viewModel = ItemListViewModel(GetItemsUseCase(fakeRepository))
    viewModel.state.test {        assertEquals(ItemListState(), awaitItem())  // initial
        viewModel.onSearch("query")        val loading = awaitItem()        assertTrue(loading.isLoading)
        val loaded = awaitItem()        assertFalse(loaded.isLoading)        assertEquals(1, loaded.items.size)    }}

Testing with TestDispatcher

kotlin
@Testfun `parallel load completes correctly`() = runTest {    val viewModel = DashboardViewModel(        itemRepo = FakeItemRepo(),        statsRepo = FakeStatsRepo()    )
    viewModel.load()    advanceUntilIdle()
    val state = viewModel.state.value    assertNotNull(state.items)    assertNotNull(state.stats)}

Faking Flows

kotlin
class FakeItemRepository : ItemRepository {    private val _items = MutableStateFlow<List<Item>>(emptyList())
    override fun observeItems(): Flow<List<Item>> = _items
    fun emit(items: List<Item>) { _items.value = items }
    override suspend fun getItemsByCategory(category: String): Result<List<Item>> {        return Result.success(_items.value.filter { it.category == category })    }}

Anti-Patterns to Avoid

  • Using GlobalScope — leaks coroutines, no structured cancellation
  • Collecting Flows in init {} without a scope — use viewModelScope.launch
  • Using MutableStateFlow with mutable collections — always use immutable copies: _state.update { it.copy(list = it.list + newItem) }
  • Catching CancellationException — let it propagate for proper cancellation
  • Using flowOn(Dispatchers.Main) to collect — collection dispatcher is the caller's dispatcher
  • Creating Flow in @Composable without remember — recreates the flow every recomposition

References

See skill: compose-multiplatform-patterns for UI consumption of Flows. See skill: android-clean-architecture for where coroutines fit in layers.

來源與署名

來源:affaan-m/ECC位於pi/core/skills/kotlin-coroutines-flows提交ef648e0

授權條款: 無授權條款

內容歸原作者所有。SourceWeft 從公開儲存庫中收錄這些內容。

檢舉或申請下架