Общее изменяемое состояние и конкурентность
Корутины можно выполнять параллельно, используя многопоточный диспетчер, например 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 не помогут
Распространено заблуждение, что объявление переменной 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", потому что 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
Это самое быстрое решение для данной задачи. Оно подходит для обычных счётчиков, коллекций, очередей и других стандартных структур данных, а также для базовых операций над ними. Однако его непросто масштабировать на сложное состояние или сложные операции, для которых нет готовых потокобезопасных реализаций.
Ограничение потока: мелкозернистый подход
Ограничение потока — это подход к проблеме общего изменяемого состояния, при котором весь доступ к определённому общему состоянию ограничен одним потоком. Его обычно используют в приложениях с пользовательским интерфейсом, где всё состояние UI закреплено за одним потоком обработки событий/приложения. Этот подход легко применять в корутинах, используя однопоточный контекст.
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
В этом примере блокировка применяется мелкозернистым образом, что снижает производительность. Однако в некоторых случаях это хороший выбор: например, когда необходимо периодически изменять общее состояние, но нет естественного потока, за которым можно было бы его закрепить.
© 2010–2026 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