DimaSOS

Корутины

Корутины и Flow: structured concurrency, отмена, обработка ошибок и то, чем StateFlow отличается от SharedFlow.

#Основы

suspend fun loadUser(id: String): User {
    return api.getUser(id)   // приостанавливает корутину, не блокирует поток
}

suspend означает «может приостановиться». Поток при этом освобождается и берёт другую работу — в отличие от блокирующего вызова.

val job = scope.launch {
    loadUser("42")           // fire-and-forget, результат не нужен
}

val deferred = scope.async {
    loadUser("42")           // нужен результат
}
val user = deferred.await()

launch возвращает Job, asyncDeferred с результатом.

// параллельно: суммарное время = время самого долгого
coroutineScope {
    val user = async { api.getUser(id) }
    val orders = async { api.getOrders(id) }
    Screen(user.await(), orders.await())
}

Два запроса одновременно. Если await() вызвать сразу после async, параллельности не будет.

Частая ошибка: val a = async { }.await() в одну строку — это последовательное выполнение.

job.cancel()
job.join()              // дождаться завершения
job.cancelAndJoin()     // отменить и дождаться

Управление жизнью корутины.

#Диспетчеры

Dispatchers.Main        // UI-поток (Android)
Dispatchers.IO          // сеть, диск, БД — пул с большим лимитом
Dispatchers.Default     // CPU-bound: парсинг, сортировка, шифрование
Dispatchers.Unconfined  // почти всегда не то, что вам нужно

IO рассчитан на блокирующие ожидания, Default ограничен числом ядер.

suspend fun parse(json: String): List<Item> =
    withContext(Dispatchers.Default) {
        Json.decodeFromString(json)
    }

Правило: suspend-функция сама отвечает за свой диспетчер. Вызывающий не должен думать, с какого потока её звать — это называется main-safety.

Поэтому в репозитории withContext внутри функции, а не launch(Dispatchers.IO) на стороне ViewModel.

// внедрять диспетчеры, а не хардкодить — иначе не подменишь в тестах
class UserRepository(
    private val api: Api,
    private val io: CoroutineDispatcher = Dispatchers.IO,
) {
    suspend fun users() = withContext(io) { api.users() }
}

Диспетчер как зависимость. В тесте подставляется StandardTestDispatcher.

#Structured concurrency

Каждая корутина принадлежит scope. Отмена scope отменяет всех детей, падение ребёнка по умолчанию валит родителя. Это не ограничение, а то, что избавляет от утечек.

coroutineScope {
    launch { a() }
    launch { b() }
}   // вернётся только когда завершатся оба
    // если один упал — отменяются все, исключение уходит наружу

coroutineScope — все или никто. Для операций, которые бессмысленны по частям.

supervisorScope {
    launch { mayFail() }    // падение не тронет соседа
    launch { other() }
}

supervisorScope изолирует падения детей друг от друга. Для независимых задач: например, три виджета на дашборде.

// НЕЛЬЗЯ в продакшене
GlobalScope.launch { upload() }

GlobalScope не отменяется ничем и живёт до смерти процесса. Гарантированная утечка и обращения к мёртвому UI.

Если работа должна переживать экран — это WorkManager или scope уровня приложения, а не GlobalScope.

class MyViewModel : ViewModel() {
    fun load() = viewModelScope.launch { /* отменится в onCleared */ }
}

// в Activity/Fragment
lifecycleScope.launch { /* отменится при уничтожении */ }

Готовые scope в Android. Своими руками scope создавать нужно редко.

#Отмена и таймауты

withTimeout(5_000) { api.slowCall() }          // бросит TimeoutCancellationException
withTimeoutOrNull(5_000) { api.slowCall() }     // вернёт null

Таймаут на операцию.

// отмена кооперативна: цикл без suspend-точек не прервётся
while (isActive) {
    doChunk()
}

// или явная проверка
ensureActive()

Тяжёлый CPU-цикл нужно проверять на отмену вручную. suspend-вызовы (delay, сетевые) проверяют сами.

try {
    doWork()
} catch (e: CancellationException) {
    throw e            // ОБЯЗАТЕЛЬНО пробросить дальше
} catch (e: Exception) {
    log(e)
}

CancellationException — механизм отмены, а не ошибка. Проглотив её, вы ломаете structured concurrency: родитель считает корутину живой.

То же касается catch (e: Throwable) и runCatching — последний ловит CancellationException тоже.

withContext(NonCancellable) {
    db.commitTransaction()   // критичная часть, доводим до конца
}

Точечно, только для очистки и коммитов. Не как способ «чтобы не отменялось».

#Ошибки

// launch: исключение уходит вверх сразу
scope.launch {
    try { risky() } catch (e: IOException) { showError(e) }
}

// async: исключение всплывает в момент await()
val d = scope.async { risky() }
try { d.await() } catch (e: IOException) { showError(e) }

Ключевая разница: у async исключение «ждёт» вызова await.

val handler = CoroutineExceptionHandler { _, e -> log(e) }
val scope = CoroutineScope(SupervisorJob() + Dispatchers.Main + handler)

CoroutineExceptionHandler работает только для корневых launch. Для async и для вложенных корутин он не сработает.

sealed interface Result<out T> {
    data class Ok<T>(val value: T) : Result<T>
    data class Err(val cause: Throwable) : Result<Nothing>
}

suspend fun users(): Result<List<User>> =
    try { Result.Ok(api.users()) }
    catch (e: CancellationException) { throw e }
    catch (e: Exception) { Result.Err(e) }

Ошибки как значение вместо исключений через все слои — удобнее для UI-состояния и не теряет отмену.

#Flow

fun ticker(): Flow<Int> = flow {
    var i = 0
    while (true) { emit(i++); delay(1_000) }
}

Cold flow: тело не выполняется, пока нет подписчика, и запускается заново для каждого.

repo.users()
    .map { list -> list.filter { it.isActive } }
    .flowOn(Dispatchers.Default)     // влияет на всё ВЫШЕ по цепочке
    .catch { e -> emit(emptyList()) }
    .onEach { render(it) }
    .launchIn(viewModelScope)

flowOn меняет контекст для upstream. catch ловит исключения только из того, что выше него — порядок операторов важен.

combine(userFlow, settingsFlow) { user, settings -> UiState(user, settings) }
zip(a, b) { x, y -> x to y }            // ждёт пару, combine — по любому обновлению
flatMapLatest { query -> search(query) }  // отменяет предыдущий поиск
debounce(300)                             // для поля ввода
distinctUntilChanged()

Рабочий набор операторов. Связка debounce + flatMapLatest — канонический живой поиск.

flow { emit(load()) }
    .retryWhen { cause, attempt ->
        if (cause is IOException && attempt < 3) { delay(1000 * (attempt + 1)); true }
        else false
    }

Повтор с нарастающей задержкой.

#StateFlow и SharedFlow

private val _state = MutableStateFlow(UiState())
val state: StateFlow<UiState> = _state.asStateFlow()

_state.update { it.copy(isLoading = true) }

StateFlow: всегда есть текущее значение, новый подписчик сразу его получает, одинаковые подряд значения не эмитятся (distinctUntilChanged встроен).

update вместо value = — атомарно, без гонки при параллельных изменениях.

private val _events = MutableSharedFlow<UiEvent>(replay = 0)
val events = _events.asSharedFlow()

SharedFlow: нет текущего значения, настраиваемый replay, повторы не фильтруются. Для событий.

val users: StateFlow<List<User>> = repo.usersFlow()
    .stateIn(
        scope = viewModelScope,
        started = SharingStarted.WhileSubscribed(5_000),
        initialValue = emptyList(),
    )

stateIn превращает cold flow в hot StateFlow. WhileSubscribed(5000) — держит подписку 5 секунд после ухода последнего подписчика, чтобы поворот экрана не перезапускал загрузку.

lifecycleScope.launch {
    repeatOnLifecycle(Lifecycle.State.STARTED) {
        viewModel.state.collect { render(it) }
    }
}

Правильный сбор во View-системе: подписка живёт только пока экран видим. В Compose эквивалент — collectAsStateWithLifecycle.

StateFlowSharedFlowChannel
Текущее значениеестьнетнет
Повторы одинаковыхфильтруютсяпроходятпроходят
Новый подписчикполучает текущеепо replayне получает
Несколько подписчиковвсе получаютвсе получаютделят события
Для чегосостояние экранасобытия для всеходнократные события

#Тестирование

@Test
fun loadsUsers() = runTest {
    val vm = UsersViewModel(FakeRepo())
    vm.load()
    advanceUntilIdle()
    assertEquals(2, vm.state.value.users.size)
}

runTest подменяет время: delay проматывается мгновенно, тест не ждёт реальных секунд.

@get:Rule val dispatcherRule = MainDispatcherRule()

class MainDispatcherRule(
    private val dispatcher: TestDispatcher = UnconfinedTestDispatcher(),
) : TestWatcher() {
    override fun starting(d: Description) = Dispatchers.setMain(dispatcher)
    override fun finished(d: Description) = Dispatchers.resetMain()
}

Без подмены Dispatchers.Main любой тест с viewModelScope упадёт с Module with the Main dispatcher had failed to initialize.

// StandardTestDispatcher — корутина стартует только на advance*
// UnconfinedTestDispatcher — стартует немедленно
advanceUntilIdle()          // выполнить всё запланированное
advanceTimeBy(1_000)
runCurrent()

Выбор диспетчера меняет момент запуска. Unconfined проще для ViewModel-тестов, Standard точнее для проверки порядка.

@Test
fun emitsStates() = runTest {
    vm.state.test {                 // app.cash.turbine
        assertFalse(awaitItem().isLoading)
        vm.load()
        assertTrue(awaitItem().isLoading)
        assertEquals(2, awaitItem().users.size)
        cancelAndIgnoreRemainingEvents()
    }
}

Turbine — стандарт для проверки последовательности значений потока. Без неё приходится вручную собирать в список.

СимптомПричина
Main dispatcher had failed to initializeне подменён Dispatchers.Main
Тест зависаетhot flow без cancel, или собирается бесконечный поток
Корутина не отменяетсяCPU-цикл без isActive
Родитель падает из-за одного ребёнканужен supervisorScope
Отмена «не работает»проглочен CancellationException (в том числе runCatching)
Загрузка перезапускается при поворотенет WhileSubscribed в stateIn
Поток собирается в фоне и жжёт батареюcollectAsState вместо collectAsStateWithLifecycle