Ранее мы разобрались c пробематикой, которая привела к созданию механизма Kotlin Coroutines а так же с теми идеями, что были заложены при реализации.

Диспатчеры

Cмысл диспатчера простой — корутина должна выполняться на каком‑то потоке. Диспатчер — это и есть объект, представляющий конкретный пул потоков.

Как и остальные объекты, определяющие поведение корутин, объект к диспатчеру можно получить доступ из контекста корутины.

val currentDispatcher = coroutineContext[ContinuationInterceptor] as CoroutineDispatcher

Указать диспатчер, на котором будет выполняться корутина можно явно при создании корутины или при определении suspend‑функции:

launc (dispatcher) {
   ... 
}

suspend fun helloWorld() = withContext(dispatcher) {
  ....
}

Если диспатчер явно не задается — диспатчер будет наследоваться от родительской корутины.

Из коробки coroutine runtime предоставляет следующие диспатчеры:

Dispatchers.Default

  • Предназначен для выполнения операций, требующих высокой нагрузки на процессор (CPU). 

  • Размер пула потоков соответствует количеству ядер на устройстве. 

Dispatchers.Main

  • Для Android запускает корутины в основном потоке (UI thread). 

  • Важно избегать блокировки этого потока. 

  • Не существует в юнит‑тестах (при необходимости можно создать собственный Main‑диспетчер).

Dispatchers.IO

  • Предназначен для выполнения блокирующих операций (ввода‑вывода, чтение/запись файлов, доступ к Shared Preferences и так далее). 

  • Размер пула потоков составляет 64 (или соответствует числу ядер, если их больше 64). 

  • Применяется для функций, выполняющих блокирующие операции. 

Dispatchers.Unconfined

  • Корутина запускается в том же потоке, в котором была запущена; смена потока может произойти после вызова вложенной корутины из дочерней. 

  • Полезен для юнит‑тестов. 

  • Не рекомендуется использовать

Диспатчер можно создать самостоятельно:

val executor = Executors.newFixedThreadPool(2)
val customDispatcher = executor.asCoroutineDispatcher()

однако делать так не рекомендуется, так как потребуется управлять жизненным циклом диспатчера в ручную.

Рекомендуемым решением является выделение на базе Dispatchers.Default или Dispatchers.IO нового пула нужного размера с помощью вызова метода limitedParallelism

val confined = Dispatchers.Default.limitedParallelism(1, "incrementDispatcher")
var counter = 0

// Invoked from arbitrary coroutines
launch(confined) {
    // This increment is sequential and race-free
    ++counter
}

Нюансы работы c корутинами, про которые не стоит забывать

ThreadLocal — переменные

Вызов корутины — это потенциальное место переключения потока выполнения. После вызова suspend‑метода поток, выполняющий корутину в общем случае меняется.

suspend fun helloWorld()  {
    val threadLocal = ThreadLocal<String>()
    threadLocal.set("main")
    println("thread local value: '${threadLocal.get()}'")
    delay(2000)
    // Может быть как "main", так и null
    println("thread local value: '${threadLocal.get()}'")
}

Если в логике приложения все же имеется необходимость использования переменных, связанных с потоком, стоит воспользоваться расширением asContextElement и добавить нужный элемент к контексту корутины:

 suspend fun helloWorld() = withContext(Dispatchers.Default) {
     val threadLocal = ThreadLocal<String>()
     threadLocal.set("main")
     launch(threadLocal.asContextElement()) {
         println("thread local value: '${threadLocal.get()}'")
         delay(2000)
         println("thread local value: '${threadLocal.get()}'")
     }
 }

Прерывание корутины

При вызове метода job.cancel() происходят две вещи:

  1. Внутренний флаг состояния Job переводится в статус Cancelling.

  2. Job проходит по списку своих детей и рекурсивно вызывает cancel() у каждого из них.

На этом работа метода cancel() заканчивается. Метод не останавливает код напрямую и сам по себе не бросает исключений в том месте, где выполняется корутина.

За остановку корутин отвечает kotlin runtime. При этом сама отмена происходит не сразу, а в ближайшей точке вызова корутины (suspension point). Если выполняется долгая синхронная операция не вызывающая других корутин — корутина может просто подвиснуть. Для решения проблемы при выполнении долгих операций нужно время от времени проверять статус корутины через вызов ensureActive().

while (true) {
    ensureActive()  // suspend-функция; проверяем не было ли отмены корутины  
    heavyOperationPart() // синхронная функция
}

CancellationException

Для сигнализации корутине, того что она была отменена kotlin runtime использует CancellationException. При этом исключение вылетит только после точки вызова корутины (suspension point). Если таковой точки нет — то и исключение не будет получено.

CancellationException — обычное исключение, которое можно обработать в блоке catch. Но CancellationException нужно пробрасывать дальше.

while (true) {
    try {
        ensureActive()  // suspend-функция; проверяем не было ли отмены корутины  
        heavyOperationPart() // синхронная функция
    } catch(e: CancellationException) {
        releaseResources()  // Освобождаем ресурсы
        // Обязательно перебрасываем отмену
        throw e
    } catch (e: Exception) {
        log.error("что-то упало", e)
    }
}

Обработка ошибок

Try.. catch

Рассмотрим простой пример:

class WorkerInvoker {
    private val realWorker = RealWorker()
    suspend fun startWorks() {
        realWorker.doWork()
    }
}

class RealWorker {
    suspend fun doWork() = withContext(Dispatchers.Default) {
        launch {
            delay(Duration.ofSeconds(2))
        }
    }
}

suspend fun main() {
    val workerInvoker = WorkerInvoker()
    workerInvoker.startWorks()
}

Несмотря на то что в методе RealWorker.doWork launch запускет корутину без блокировки текущего потока выполнения, благодаря механизму Structured Concurrency сначала завершится RealWorker.doWork, потом WorkerInvoker.startWorks и только потом функция main.

Теперь модифицируем пример:

class RealWorker {
    suspend fun doWork() = withContext(Dispatchers.Default) {
        launch {
            delay(Duration.ofSeconds(2))
            throw Exception("doWork Exception")
        }
    }
}

Получим такой трейс:

Exception in thread "main" java.lang.Exception: DoWork Exception
	at ru.voskhod.RealWorker$doWork$2$1.invokeSuspend(Main.kt:23)
	at kotlin.coroutines.jvm.internal.BaseContinuationImpl.resumeWith(ContinuationImpl.kt:34)
  ....
    

Модифицируем код:

suspend fun doWork() = withContext(Dispatchers.Default) {
        try {
            launch {
                delay(Duration.ofSeconds(2))
                throw Exception("DoWork Exception")
            }
        } catch (e: Exception) {
            e.printStackTrace()
            throw Exception("RealWorker Exception")
        }
    }

Трейс не меняется:

Exception in thread "main" java.lang.Exception: DoWork Exception
	at ru.voskhod.RealWorker$doWork$2$1.invokeSuspend(Main.kt:24)
	at kotlin.coroutines.jvm.internal.BaseContinuationImpl.resumeWith(ContinuationImpl.kt:34)
  ....

Модифицируем код ещё раз:


class WorkerInvoker {
    private val realWorker = RealWorker()
    suspend fun startWorks() {
        try {
            realWorker.doWork()
        } catch (e: Exception) {
            e.printStackTrace()
            throw Exception("WorkerInvoker Exception")
        }
    }
}

class RealWorker {
    suspend fun doWork() = withContext(Dispatchers.Default) {
        launch {
            delay(Duration.ofSeconds(2))
            throw Exception("DoWork Exception")
        }
    }
}

В выводе:

java.lang.Exception: DoWork Exception
	at ru.voskhod.RealWorker$doWork$2$1.invokeSuspend(Main.kt:29)
	at kotlin.coroutines.jvm.internal.BaseContinuationImpl.resumeWith(ContinuationImpl.kt:34)
	r.kt:704)
      ....
Exception in thread "main" java.lang.Exception: WorkerInvoker Exception
	at ru.voskhod.WorkerInvoker.doWork(Main.kt:20)
	at ru.voskhod.WorkerInvoker$doWork$1.invokeSuspend(Main.kt)
	at kotlin.coroutines.jvm.internal.BaseContinuationImpl.resumeWith(ContinuationImpl.kt:34)
	at kotlinx.coroutines.internal.DispatchedContinuationKt.resumeCancellableWithInternal(DispatchedContinuation.kt:278)
	at kotlinx.coroutines.DispatchedCoroutine.afterResume(Builders.common.kt:588)
	  ....

Почему же с try.. catch в методе RealWorker.doWork не отработал, а отработал только в методе WorkerInvoker.doWork?

Дело в том, что launch запускает корутину асинхронно и к моменту завершения launch‑корутины код блока RealWorker.doWork уже выполнен. Корутина, в которой выполняется RealWorker.doWork просто ждет завершения выполнения дочерних корутин. В случае же с WorkerInvoker.doWork в выполнение метода приостанавливается до завершения вызова realWorker.doWork().

CoroutineExceptionHandler

Для обработки исключений, возникающий в дочерних корутинах существует опциональный для контекста объект CoroutineExceptionHandler.

Есть два моменты, связанных с CoroutineExceptionHandler:

  • CoroutineExceptionHandler вызывается внутри kotlin runtime и поток в котором он вызывается не определяется. То есть обработчик должен быть потокобезопасным и быстро завершаться

  • CoroutineExceptionHandler сработает только если его установить в корутине верхнего уровня (корневой корутине). Если его установить в дочерней корутине — он не сработает:

// Так делять нельзя!!
class RealWorker {
    suspend fun doWork() = withContext(Dispatchers.Default) {
        withContext(CoroutineExceptionHandler { ctx, ex ->
            println("Exception $ex thrown from coroutine context $ctx")
        }) {
            launch {
                delay(Duration.ofSeconds(2))
                throw Exception("DoWork Exception1")
            }
        }
    }
}

Но выходе все так же:

Exception in thread "main" java.lang.Exception: DoWork Exception1
	at ru.voskhod.RealWorker$doWork$2$2$1.invokeSuspend(Main.kt:34)
	at kotlin.coroutines.jvm.internal.BaseContinuationImpl.resumeWith(ContinuationImpl.kt:34)
	.kt:704)
    .....

Для использования CoroutineExceptionHandler надо создавать новую иерархию корутин:

class RealWorker {
    suspend fun doWork() = withContext(Dispatchers.Default) {

        val handler = CoroutineExceptionHandler { ctx, ex ->
            println("Exception $ex thrown from coroutine context $ctx")
        }
        val scope = CoroutineScope(SupervisorJob() + handler)

        val job = scope.launch {
            delay(Duration.ofSeconds(2))
            throw Exception("DoWork Exception1")
        }
        job.join() // Иерархия новая, нужно явно дожидаться выполнения job
    }
}

Программа завершилась успешно и вывела:

> Task :org.test.sample.MainKt.main()
Exception java.lang.Exception: DoWork Exception1 thrown from coroutine context [org.test.sample.RealWorker$doWork$2$invokeSuspend$$inlined$CoroutineExceptionHandler$1@7f05f394, StandaloneCoroutine{Cancelling}@c68c8d1, Dispatchers.Default]

Напомню, что при возникновении исключения, отменяется вся иерархия корутин. Если не хочется прерывать остальные дочерние корутины или нужно получить все исключения стоит использовать CoroutineExceptionHandler в связке с SupervisorJob()/supervisorScope:

suspend fun doWork() = withContext(Dispatchers.Default) {

        val handler = CoroutineExceptionHandler { ctx, ex ->
            println("Exception $ex thrown from coroutine context $ctx")
        }

        val scope = CoroutineScope(SupervisorJob() + handler)

        val jobList = with(scope) {
            listOf(
                launch {
                    supervisorScope {
                        launch {
                            delay(Duration.ofSeconds(2))
                            throw Exception("DoWork Exception1")
                        }
                        launch {
                            delay(Duration.ofSeconds(3))
                            throw Exception("DoWork Exception2")
                        }
                    }
                },
                launch {
                    delay(Duration.ofMillis(400))
                    throw Exception("DoWork Exception3")
                }
            )
        }

        jobList.joinAll()
    }

Программа выведет все исключения:

> Task :org.test.sample.MainKt.main()
Exception java.lang.Exception: DoWork Exception3 thrown from coroutine context [org.test.sample.RealWorker$doWork$2$invokeSuspend$$inlined$CoroutineExceptionHandler$1@64f0b65b, StandaloneCoroutine{Cancelling}@2dcedeb, Dispatchers.Default]
Exception java.lang.Exception: DoWork Exception1 thrown from coroutine context [org.test.sample.RealWorker$doWork$2$invokeSuspend$$inlined$CoroutineExceptionHandler$1@64f0b65b, StandaloneCoroutine{Cancelling}@3444aebc, Dispatchers.Default]
Exception java.lang.Exception: DoWork Exception2 thrown from coroutine context [org.test.sample.RealWorker$doWork$2$invokeSuspend$$inlined$CoroutineExceptionHandler$1@64f0b65b, StandaloneCoroutine{Cancelling}@29cbffc0, Dispatchers.Default]

Ссылки и благодарности