Spec-Zone.ru › Kotlin 1.4

Содержание

  • Общее изменяемое состояние и конкурентность
    • Проблема
    • Переменные volatile не помогают
    • Потокобезопасные структуры данных
    • Внедрение потоков (мелкозернистое)
    • Внедрение потоков (крупнозернистое)
    • Взаимное исключение
    • Акторы

Общее изменяемое состояние и конкурентность

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

Существует распространённое заблуждение, что объявление переменной 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" в конце, потому что переменные volatile гарантируют линейно-выполнимые (это технический термин для «атомарные») чтение и запись в соответствующую переменную, но не обеспечивают атомарность более крупных операций (инкремента в нашем случае).

Потокобезопасные структуры данных

Общим решением, которое работает как для потоков, так и для корутин, является использование потокобезопасной (также синхронизированной, линейно-выполнимой или атомарной) структуры данных, которая обеспечивает необходимую синхронизацию для соответствующих операций, которые необходимо выполнить с общим состоянием. В случае простого счётчика мы можем использовать класс 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's sealed-классы хорошо подходят для этой цели. Мы определяем CounterMsg sealed-класс с 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    

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

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

END_OF_DOCUMENT_MARKER

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

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

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

Spec-Zone.ru

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