Spec-Zone.ru › Kotlin 1.6

Управляемое состояние и конкурентность

Корутины могут выполняться параллельно с помощью многопоточного диспетчера, такого как Dispatchers.Default. Это создаёт все обычные проблемы параллелизма. Основной проблемой является синхронизация доступа к управляемому общему состоянию. Некоторые решения этой проблемы в мире корутин похожи на решения в многопоточной среде, но другие уникальны.

Проблема

Давайте запустим сто корутин, каждая из которых выполняет одно и то же действие тысячу раз. Мы также измерим время их завершения для дальнейшего сравнения:

suspend fun massiveRun(action: suspend () -> Unit) {
    val n = 100  // number of coroutines to launch
    val k = 1000 // times an action is repeated by each coroutine
    val time = measureTimeMillis {
        coroutineScope { // scope for coroutines 
            repeat(n) {
                launch {
                    repeat(k) { action() }
                }
            }
        }
    }
    println("Completed ${n * k} actions in $time ms")    
}

Начнём с очень простого действия, которое увеличивает общую изменяемую переменную с использованием многопоточного Dispatchers.Default.

import kotlinx.coroutines.*
import kotlin.system.*    

suspend fun massiveRun(action: suspend () -> Unit) {
    val n = 100  // number of coroutines to launch
    val k = 1000 // times an action is repeated by each coroutine
    val time = measureTimeMillis {
        coroutineScope { // scope for coroutines 
            repeat(n) {
                launch {
                    repeat(k) { action() }
                }
            }
        }
    }
    println("Completed ${n * k} actions in $time ms")    
}

//sampleStart
var counter = 0

fun main() = runBlocking {
    withContext(Dispatchers.Default) {
        massiveRun {
            counter++
        }
    }
    println("Counter = $counter")
}
//sampleEnd    

Вы можете получить полный код здесь.

Что будет выведено в конце? Вряд ли когда-либо будет напечатано «Counter = 100000», потому что сто корутин увеличивают counter одновременно из нескольких потоков без какой-либо синхронизации.

Изменчивые переменные не помогают

Существует распространённое заблуждение, что сделать переменную volatile решает проблему конкурентности. Давайте попробуем:

import kotlinx.coroutines.*
import kotlin.system.*

suspend fun massiveRun(action: suspend () -> Unit) {
    val n = 100  // number of coroutines to launch
    val k = 1000 // times an action is repeated by each coroutine
    val time = measureTimeMillis {
        coroutineScope { // scope for coroutines 
            repeat(n) {
                launch {
                    repeat(k) { action() }
                }
            }
        }
    }
    println("Completed ${n * k} actions in $time ms")    
}

//sampleStart
@Volatile // in Kotlin `volatile` is an annotation 
var counter = 0

fun main() = runBlocking {
    withContext(Dispatchers.Default) {
        massiveRun {
            counter++
        }
    }
    println("Counter = $counter")
}
//sampleEnd    

Вы можете получить полный код здесь.

Этот код работает медленнее, но мы всё ещё не получаем «Counter = 100000» в конце, потому что изменчивые переменные гарантируют линейно упорядоченные (это технический термин для «атомарных») чтение и запись в соответствующую переменную, но не обеспечивают атомарность более крупных действий (в нашем случае – инкремент).

Безопасные для потоков структуры данных

Общим решением, которое работает как для потоков, так и для корутин, является использование структуры данных, безопасной для потоков (также известной как синхронизированной, линейно упорядоченной или атомарной), которая обеспечивает необходимую синхронизацию для соответствующих операций, которые необходимо выполнить над общим состоянием. В случае простого счётчика мы можем использовать AtomicInteger класс, который имеет атомарные incrementAndGet операции:

import kotlinx.coroutines.*
import java.util.concurrent.atomic.*
import kotlin.system.*

suspend fun massiveRun(action: suspend () -> Unit) {
    val n = 100  // number of coroutines to launch
    val k = 1000 // times an action is repeated by each coroutine
    val time = measureTimeMillis {
        coroutineScope { // scope for coroutines 
            repeat(n) {
                launch {
                    repeat(k) { action() }
                }
            }
        }
    }
    println("Completed ${n * k} actions in $time ms")    
}

//sampleStart
val counter = AtomicInteger()

fun main() = runBlocking {
    withContext(Dispatchers.Default) {
        massiveRun {
            counter.incrementAndGet()
        }
    }
    println("Counter = $counter")
}
//sampleEnd    

Вы можете получить полный код здесь.

Это самое быстрое решение для этой конкретной проблемы. Оно работает для обычных счётчиков, коллекций, очередей и других стандартных структур данных и основных операций над ними. Однако это не всегда легко масштабировать для сложного состояния или сложных операций, для которых нет готовых реализаций безопасных для потоков структур.

Файное зернистое ограничение потока

Ограничение потока — подход к проблеме общего изменяемого состояния, где весь доступ к конкретному общему состоянию ограничен одним потоком. Обычно используется в приложениях для пользовательского интерфейса, где всё состояние пользовательского интерфейса ограничено одним потоком обработки событий/приложения. Легко применять с корутинами, используя однопоточную контекстную среду.

import kotlinx.coroutines.*
import kotlin.system.*

suspend fun massiveRun(action: suspend () -> Unit) {
    val n = 100  // number of coroutines to launch
    val k = 1000 // times an action is repeated by each coroutine
    val time = measureTimeMillis {
        coroutineScope { // scope for coroutines 
            repeat(n) {
                launch {
                    repeat(k) { action() }
                }
            }
        }
    }
    println("Completed ${n * k} actions in $time ms")    
}

//sampleStart
val counterContext = newSingleThreadContext("CounterContext")
var counter = 0

fun main() = runBlocking {
    withContext(Dispatchers.Default) {
        massiveRun {
            // confine each increment to a single-threaded context
            withContext(counterContext) {
                counter++
            }
        }
    }
    println("Counter = $counter")
}
//sampleEnd      

Вы можете получить полный код здесь.

Этот код работает очень медленно, потому что он использует тонкое ограничение потока. Каждый отдельный инкремент переключается с многопоточного Dispatchers.Default контекста на однопоточный контекст с помощью withContext(counterContext) блока.

Грубозернистое ограничение потока

На практике ограничение потока выполняется большими блоками, например, большие фрагменты логики обновления состояния ограничены одним потоком. Следующий пример делает это таким образом, запуская каждую корутину в однопоточном контексте вначале.

import kotlinx.coroutines.*
import kotlin.system.*

suspend fun massiveRun(action: suspend () -> Unit) {
    val n = 100  // number of coroutines to launch
    val k = 1000 // times an action is repeated by each coroutine
    val time = measureTimeMillis {
        coroutineScope { // scope for coroutines 
            repeat(n) {
                launch {
                    repeat(k) { action() }
                }
            }
        }
    }
    println("Completed ${n * k} actions in $time ms")    
}

//sampleStart
val counterContext = newSingleThreadContext("CounterContext")
var counter = 0

fun main() = runBlocking {
    // confine everything to a single-threaded context
    withContext(counterContext) {
        massiveRun {
            counter++
        }
    }
    println("Counter = $counter")
}
//sampleEnd     

Вы можете получить полный код здесь.

Теперь это работает гораздо быстрее и даёт правильный результат.

Взаимное исключение

Решение проблемы взаимного исключения заключается в защите всех изменений общего состояния с помощью критической секции, которая никогда не выполняется одновременно. В блокирующем мире вы обычно использовали бы synchronized или ReentrantLock для этого. Альтернативой корутин является Mutex. Он имеет функции lock и unlock для определения критической секции. Ключевое различие состоит в том, что Mutex.lock() — это приостанавливаемая функция. Она не блокирует поток.

Также есть расширение функции withLock, которая удобно представляет mutex.lock(); try { ... } finally { mutex.unlock() } шаблон:

import kotlinx.coroutines.*
import kotlinx.coroutines.sync.*
import kotlin.system.*

suspend fun massiveRun(action: suspend () -> Unit) {
    val n = 100  // number of coroutines to launch
    val k = 1000 // times an action is repeated by each coroutine
    val time = measureTimeMillis {
        coroutineScope { // scope for coroutines 
            repeat(n) {
                launch {
                    repeat(k) { action() }
                }
            }
        }
    }
    println("Completed ${n * k} actions in $time ms")    
}

//sampleStart
val mutex = Mutex()
var counter = 0

fun main() = runBlocking {
    withContext(Dispatchers.Default) {
        massiveRun {
            // protect each increment with lock
            mutex.withLock {
                counter++
            }
        }
    }
    println("Counter = $counter")
}
//sampleEnd    

Вы можете получить полный код здесь.

Блокировка в этом примере является тонкой, поэтому она несёт соответствующие затраты. Однако это хороший выбор в некоторых ситуациях, когда вы абсолютно должны периодически изменять некоторое общее состояние, но нет естественного потока, к которому это состояние ограничено.

Акторы

Актер — это сущность, состоящая из сочетания корутины, состояния, ограниченного и инкапсулированного в этой корутине, и канала для связи с другими корутинами. Простой актер может быть написан как функция, но актер со сложным состоянием лучше подходит для класса.

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

Первый шаг использования актера — определение класса сообщений, которые будет обрабатывать актер. Для этой цели хорошо подходят запечатанные классы Kotlin. Определяем CounterMsg запечатанный класс с IncCounter сообщением для инкремента счётчика и GetCounter сообщением для получения его значения. Последнее должно отправить ответ. Для этой цели используется примитив обмена CompletableDeferred, представляющий собой единственное значение, которое будет известно (передано) в будущем.

// Message types for counterActor
sealed class CounterMsg
object IncCounter : CounterMsg() // one-way message to increment counter
class GetCounter(val response: CompletableDeferred<Int>) : CounterMsg() // a request with reply

Затем мы определяем функцию, которая запускает актера с помощью билдера корутин actor:

// This function launches a new counter actor
fun CoroutineScope.counterActor() = actor<CounterMsg> {
    var counter = 0 // actor state
    for (msg in channel) { // iterate over incoming messages
        when (msg) {
            is IncCounter -> counter++
            is GetCounter -> msg.response.complete(counter)
        }
    }
}

Основной код прост:

import kotlinx.coroutines.*
import kotlinx.coroutines.channels.*
import kotlin.system.*

suspend fun massiveRun(action: suspend () -> Unit) {
    val n = 100  // number of coroutines to launch
    val k = 1000 // times an action is repeated by each coroutine
    val time = measureTimeMillis {
        coroutineScope { // scope for coroutines 
            repeat(n) {
                launch {
                    repeat(k) { action() }
                }
            }
        }
    }
    println("Completed ${n * k} actions in $time ms")    
}

// Message types for counterActor
sealed class CounterMsg
object IncCounter : CounterMsg() // one-way message to increment counter
class GetCounter(val response: CompletableDeferred<Int>) : CounterMsg() // a request with reply

// This function launches a new counter actor
fun CoroutineScope.counterActor() = actor<CounterMsg> {
    var counter = 0 // actor state
    for (msg in channel) { // iterate over incoming messages
        when (msg) {
            is IncCounter -> counter++
            is GetCounter -> msg.response.complete(counter)
        }
    }
}

//sampleStart
fun main() = runBlocking<Unit> {
    val counter = counterActor() // create the actor
    withContext(Dispatchers.Default) {
        massiveRun {
            counter.send(IncCounter)
        }
    }
    // send a message to get a counter value from an actor
    val response = CompletableDeferred<Int>()
    counter.send(GetCounter(response))
    println("Counter = ${response.await()}")
    counter.close() // shutdown the actor
}
//sampleEnd    

Полный код можно найти здесь.

Для корректности не имеет значения, в каком контексте выполняется сам актер. Актер — это корутина, а корутина выполняется последовательно, поэтому ограничение состояния конкретной корутиной является решением проблемы разделяемого изменяемого состояния. Действительно, актеры могут изменять собственное частное состояние, но могут влиять друг на друга только через сообщения (избегая необходимости каких-либо блокировок).

Актер более эффективен, чем блокировка под нагрузкой, потому что в этом случае у него всегда есть работа, и ему не нужно переключаться на другой контекст.

Обратите внимание, что билдер корутин actor — это дуаль билдера корутин produce. Актер связан с каналом, из которого получает сообщения, в то время как производитель связан с каналом, в который отправляет элементы.

Последнее изменение: 04 апреля 2022
Обработка исключений корутин Выражение select (экспериментальное)

© 2010–2022 JetBrains s.r.o. and Kotlin Programming Language contributors
Licensed under the Apache License, Version 2.0.
https://kotlinlang.org/docs/shared-mutable-state-and-concurrency.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API