Корутины
Корутины и 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, async — Deferred с результатом.
// параллельно: суммарное время = время самого долгого
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.
| StateFlow | SharedFlow | Channel | |
|---|---|---|---|
| Текущее значение | есть | нет | нет |
| Повторы одинаковых | фильтруются | проходят | проходят |
| Новый подписчик | получает текущее | по 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 |