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 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() происходят две вещи:
Внутренний флаг состояния Job переводится в статус Cancelling.
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]
Ссылки и благодарности
KioskNews shows a cleaned-up reading view extracted from the publisher’s page — the original always lives on their site, not ours.