Spec-Zone.ru › Kotlin 2

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

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

Полный код можно посмотреть здесь.

В этом примере блокировка применяется мелкозернистым образом, что снижает производительность. Однако в некоторых случаях это хороший выбор: например, когда необходимо периодически изменять общее состояние, но нет естественного потока, за которым можно было бы его закрепить.

27 сентября 2024 г.
Обработка исключений в корутинахВыражение select (экспериментальное)

© 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

Spec-Zone.ru

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