Kotlin Coroutines — разбираемся с диспатчерами и обработкой ошибок

—

от автора

Ранее мы разобрались 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 coroutineslaunch(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 Exceptionat 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 Exceptionat 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 Exceptionat 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 Exceptionat 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 Exception1at 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]

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

ссылка на оригинал статьи https://habr.com/ru/articles/1090314/