Общий изменяемый и конкурентный доступ к состоянию
Корутины могут выполняться параллельно с помощью многопотокового диспетчера, такого как 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
Что будет выведено в конце? Вероятность, что будет напечатано "Счётчик = 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
Этот код работает медленнее, но мы по-прежнему не получаем "Счётчик = 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–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