Лекции

Senior

Kotlin Coroutines: углублённо

Жизненный цикл Job, распространение ошибок, supervision, кооперативная отмена, dispatcher и анализ конкурентного кода.

coroutinesjobssupervisiondispatchersconcurrency
На этой странице

Это русскоязычная подготовленная версия предоставленного конспекта. Все исходные кодовые блоки сохранены как read-only примеры. Примеры не компилировались и не запускались в ходе подготовки.

Теория

Job — lifecycle coroutine

Job — не просто handle для cancel().

Основные состояния можно представить так:

           start
NEW ───────────────► ACTIVE
                       │
             ┌─────────┴─────────┐
             │                   │
          success             failure /
             │                cancellation
             ▼                   │
        COMPLETING          CANCELLING
             │                   │
             └─────────┬─────────┘
                       ▼
                   COMPLETED

Практически:

job.isActive
job.isCompleted
job.isCancelled

Официальная документация: Job API.

Почему существует Completing

val job = launch {
    launch {
        delay(5000)
    }

    println("parent body finished")
}

Тело parent завершилось, но Job ждёт child.

Parent body:
█████

Child:
██████████████████████████

Parent Job:
██████████████████████████

Завершение тела coroutine и завершение Job — не обязательно один момент.

Официальная документация: Job API.

Cancellation сверху вниз

val parent = scope.launch {
    launch { delay(10_000) }
    launch { delay(10_000) }
}

parent.cancel()
parent.cancel()
      │
      ▼
 Parent cancelled
      │
 ┌────┴────┐
 ▼         ▼
Child A   Child B
cancel    cancel

Failure child снизу вверх

scope.launch {
    launch {
        delay(100)
        error("BOOM")
    }

    launch {
        delay(10_000)
        println("finished")
    }
}

Обычный Job:

Child A failure
      ↓
    Parent
      ↓
 Child B cancellation

Failure идёт вверх, cancellation затем распространяется вниз.

Официальная документация: Job API.

CancellationException — не обычный failure

RuntimeException
       ↓
     failure
       ↓
cancel parent

CancellationException
       ↓
 cancellation

Поэтому cancellation не следует превращать в business error.

Cooperative cancellation

val job = launch(Dispatchers.Default) {
    while (true) {
        calculateSomething()
    }
}

job.cancel()

Если нет cancellable suspension points или ручных проверок, loop может продолжать работу.

Используй:

while (isActive) {
    calculate()
}

или:

while (true) {
    ensureActive()
    calculate()
}

Официальная документация: Dispatchers API.

isActive vs ensureActive()

isActive позволяет самостоятельно выйти из цикла.

ensureActive() при отмене бросает CancellationException, сохраняя стандартную cancellation semantics.

yield()

while (true) {
    processChunk()
    yield()
}

yield() проверяет cancellation и даёт scheduler возможность выполнить другие coroutine.

Но он не должен маскировать плохо спроектированную тяжёлую CPU-задачу.

finally при cancellation

try {
    doWork()
} finally {
    cleanup()
}

finally выполняется и при cancellation.

Но cancellable suspend operation внутри finally может сразу снова получить cancellation.

Официальная документация: Документация Kotlin Coroutines.

NonCancellable

Если короткий suspend-cleanup обязан завершиться:

finally {
    withContext(NonCancellable) {
        saveState()
    }
}

Не стоит помещать туда долгую business operation или огромный network sync.

Официальная документация: Документация Kotlin Coroutines.

Почему try/catch вокруг launch не ловит child failure

try {
    scope.launch {
        throw RuntimeException("Boom")
    }
} catch (e: Exception) {
    println("Caught")
}

Внешний try окружает запуск coroutine, а exception возникает позже внутри её выполнения.

Лови внутри coroutine или на корректной structured-concurrency boundary.

async и failure

val deferred = scope.async {
    error("Boom")
}

await() выбросит failure, но нельзя делать вывод, что до await() failure не влияет на hierarchy.

coroutineScope {
    val deferred = async {
        error("Boom")
    }

    delay(10_000)
    deferred.await()
}

Child async может отменить coroutineScope сразу после failure.

coroutineScope как fail-fast boundary

suspend fun loadPage(): Page =
    coroutineScope {
        val user = async { loadUser() }
        val posts = async { loadPosts() }

        Page(
            user.await(),
            posts.await(),
        )
    }

Если все части нужны для единого результата, fail-fast часто является правильной семантикой.

Официальная документация: Документация Kotlin Coroutines.

supervisorScope для независимых частей

supervisorScope {
    launch { loadProfile() }
    launch { loadAds() }
}

Падение Ads не обязано отменять Profile.

Ловушка SupervisorJob

val scope = CoroutineScope(
    SupervisorJob() + Dispatchers.Default
)

scope.launch {
    launch { taskA() }
    launch { taskB() }
}

Hierarchy:

SupervisorJob
      │
      ▼
 outer launch
      │
   ┌──┴──┐
   ▼     ▼
   A     B

A и B имеют обычного parent — outer launch.

Если A падает:

A fails
  ↓
outer launch fails
  ↓
B cancelled

SupervisorJob не превращает всё поддерево в supervision tree.

Если нужна isolation именно между A и B — используй supervisorScope на нужном уровне.

Официальная документация: CoroutineScope API.

CoroutineExceptionHandler — не business error handler

Handler работает с необработанными coroutine exceptions в соответствующей root boundary.

Для ожидаемых ошибок:

try {
    repository.load()
} catch (e: IOException) {
    // map to UI/domain error
}

Race condition

var counter = 0

coroutineScope {
    repeat(1000) {
        launch(Dispatchers.Default) {
            counter++
        }
    }
}

counter++ — read-modify-write.

Thread A          Thread B

read 10           read 10
+1                +1
write 11          write 11

Один increment потерян.

Официальная документация: Dispatchers API.

Mutex и альтернативы

val mutex = Mutex()

mutex.withLock {
    counter++
}

Но Mutex — не автоматический ответ.

Также возможны:

  • immutable state;
  • Atomic*;
  • MutableStateFlow.update;
  • Channel/actor-like serialization;
  • thread confinement.

Официальная документация: Flow API.

Thread confinement

Вместо большого количества locks можно сделать одного owner состояния:

Event A ──┐
Event B ──┼──► State owner ──► State
Event C ──┘

Если mutation сериализована, многие races исчезают архитектурно.

Dispatchers.Default

Для CPU work:

withContext(Dispatchers.Default) {
    calculateMoves()
}

Примеры:

  • legal move generation;
  • game-tree search;
  • heuristic evaluation;
  • sorting;
  • compression;
  • parsing;
  • cryptography.

Официальная документация: Dispatchers API.

Dispatchers.IO

Для blocking I/O:

withContext(Dispatchers.IO) {
    legacyBlockingApi()
}

Например blocking file/socket/JDBC APIs.

Официальная документация: Dispatchers API.

limitedParallelism

val dispatcher =
    Dispatchers.IO.limitedParallelism(4)

Можно ограничить concurrency конкретной подсистемы.

Для CPU-задачи:

val botDispatcher =
    Dispatchers.Default.limitedParallelism(2)

Официальная документация: Dispatchers API.

Dispatchers.Unconfined

Coroutine начинает выполнение в текущем call stack/thread, а после suspension может продолжиться там, где её возобновит suspend operation.

В обычном Android application code почти никогда не нужен.

Официальная документация: Документация Kotlin Coroutines.

Dispatcher ≠ новый thread

withContext(Dispatchers.IO) { ... }

не означает new Thread().

Dispatcher задаёт scheduling policy.

Официальная документация: Dispatchers API.

Thread identity не гарантируется

Coroutine может до и после suspension физически выполняться на разных threads, сохраняя тот же dispatcher/context contract.

Не строй correctness на Thread.currentThread() identity.

ThreadLocal и coroutines

Обычный ThreadLocal привязан к physical thread, а coroutine может мигрировать между threads.

Для специальных случаев существует ThreadLocal.asContextElement().

Interview case: обычный scope

coroutineScope {
    launch {
        delay(100)
        throw RuntimeException("A")
    }

    launch {
        try {
            delay(10_000)
            println("B completed")
        } finally {
            println("B finally")
        }
    }
}

Поведение:

A throws
   ↓
parent scope fails
   ↓
B cancelled
   ↓
B finally executes
   ↓
RuntimeException propagates to caller

Interview case: supervisor

supervisorScope {
    launch {
        delay(100)
        throw RuntimeException("A")
    }

    launch {
        delay(500)
        println("B completed")
    }
}

Failure A не отменяет B. Но exception A всё равно должен иметь корректную handling boundary.

Interview case: лишний вложенный launch

Плохо:

viewModelScope.launch {
    try {
        launch {
            repository.load()
        }
    } catch (e: Exception) {
        showError()
    }
}

Внешний try не ловит асинхронный child failure так, как ожидает автор.

Часто вложенный launch вообще не нужен:

viewModelScope.launch {
    try {
        repository.load()
    } catch (e: CancellationException) {
        throw e
    } catch (e: Exception) {
        showError()
    }
}

Универсальный алгоритм разбора coroutine-кода

На собеседовании:

  1. Нарисуй Job hierarchy.
  2. Определи parent каждого child.
  3. Найди supervision boundaries.
  4. Найди место failure/cancellation.
  5. Проследи failure вверх.
  6. Проследи cancellation вниз.
  7. Только затем анализируй try/catch, await и handler.

Официальная документация: Job API.

Уточнения семантики dispatcher и ThreadLocal

limitedParallelism(n) ограничивает число задач, одновременно исполняемых через этот dispatcher view; он не ограничивает общее число входящих запросов или число созданных coroutine. Задача, которая уже приостановилась, не занимает поток исполнения, поэтому этот API не является rate limiter для HTTP-запросов.

ThreadLocal.asContextElement() устанавливает значение ThreadLocal при переключении coroutine на поток. Изменение ThreadLocal внутри coroutine не становится автоматически новым значением для всех последующих выполнений: чтобы явно сменить контекстное значение, создайте новый element, например через withContext(threadLocal.asContextElement(value)).

Практика

Вопросы для разбора

Для каждого примера нарисуйте иерархию Job, определите CoroutineContext и Dispatcher, отметьте точки suspension/cancellation, затем проследите failure вверх и cancellation вниз. Для Flow дополнительно определите cold/hot semantics, владельца collection, поведение медленного consumer, replay и последствия lifecycle stop/start.

Проверка архитектурного решения

Спросите: кто владеет этой работой, кто может её отменить, где обрабатывается ошибка, нужна ли каждая emission или только latest, и переживает ли результат пересоздание UI. Coroutine не переживает process death; для гарантированной фоновой работы нужен durable scheduler, например WorkManager.

Дополнительные материалы