Разделяемое изменяемое состояние и конкурентность
Корутины могут выполняться параллельно с помощью многопоточного диспетчера, например, 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
Не имеет значения (для корректности), в каком контексте сам актор выполняется. Актор — это корутина, а корутина выполняется последовательно, поэтому ограничение состояния конкретной корутиной является решением проблемы совместного использования изменяемого состояния. Действительно, акторы могут изменять собственное частное состояние, но могут влиять друг на друга только через сообщения (избегая необходимости каких-либо блокировок).
Актор более эффективен, чем блокировка под нагрузкой, потому что в этом случае у него всегда есть работа, и он не должен переключаться на другой контекст.
© 2010–2023 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