Spec-Zone.ru › Kotlin 1.8

Разделяемое изменяемое состояние и конкурентность

Корутины могут выполняться параллельно с помощью многопоточного диспетчера, например, 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. Актор связан с каналом, из которого он получает сообщения, в то время как производитель связан с каналом, в который он отправляет элементы.

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

© 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

Spec-Zone.ru

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