Каналы
Отложенные значения предоставляют удобный способ передачи одного значения между сопрограммами. Каналы обеспечивают способ передачи потока значений.
Основы каналов
Канал Channel по своей сути очень похож на BlockingQueue. Ключевое различие заключается в том, что вместо блокирующей put операции он имеет приостанавливаемую send, а вместо блокирующей take операции — приостанавливаемую receive.
import kotlinx.coroutines.*
import kotlinx.coroutines.channels.*
fun main() = runBlocking {
//sampleStart
val channel = Channel<Int>()
launch {
// this might be heavy CPU-consuming computation or async logic, we'll just send five squares
for (x in 1..5) channel.send(x * x)
}
// here we print five received integers:
repeat(5) { println(channel.receive()) }
println("Done!")
//sampleEnd
}
Вывод этого кода:
1 4 9 16 25 Done!
Закрытие и итерация по каналам
В отличие от очереди, канал может быть закрыт, чтобы указать, что больше элементов не поступает. На стороне получателя удобно использовать обычный цикл for для получения элементов из канала.
По сути, закрытие close подобно отправке специального маркера закрытия в канал. Итерация прекращается сразу после получения этого маркера, поэтому гарантируется, что все элементы, отправленные до закрытия, будут получены:
import kotlinx.coroutines.*
import kotlinx.coroutines.channels.*
fun main() = runBlocking {
//sampleStart
val channel = Channel<Int>()
launch {
for (x in 1..5) channel.send(x * x)
channel.close() // we're done sending
}
// here we print received values using `for` loop (until the channel is closed)
for (y in channel) println(y)
println("Done!")
//sampleEnd
}
Создание производителей каналов
Шаблон, где сопрограмма производит последовательность элементов, довольно распространён. Это часть паттерна «производитель-потребитель», который часто встречается в конкурентном коде. Можно абстрагировать такого производителя в функцию, принимающую канал в качестве параметра, но это противоречит здравому смыслу, так как результаты должны возвращаться из функций.
Есть удобный конструктор сопрограмм под названием produce, который упрощает эту задачу непосредственно на стороне производителя, и функция расширения consumeEach, которая заменяет цикл for на стороне потребителя:
import kotlinx.coroutines.*
import kotlinx.coroutines.channels.*
fun CoroutineScope.produceSquares(): ReceiveChannel<Int> = produce {
for (x in 1..5) send(x * x)
}
fun main() = runBlocking {
//sampleStart
val squares = produceSquares()
squares.consumeEach { println(it) }
println("Done!")
//sampleEnd
}
Цепочки
Цепочка — это шаблон, где одна сопрограмма производит, возможно, бесконечный поток значений:
fun CoroutineScope.produceNumbers() = produce<Int> {
var x = 1
while (true) send(x++) // infinite stream of integers starting from 1
}
А другая сопрограмма или сопрограммы потребляют этот поток, выполняют некоторую обработку и производят другие результаты. В примере ниже числа просто возводятся в квадрат:
fun CoroutineScope.square(numbers: ReceiveChannel<Int>): ReceiveChannel<Int> = produce {
for (x in numbers) send(x * x)
}
Основной код запускает и соединяет всю цепочку:
import kotlinx.coroutines.*
import kotlinx.coroutines.channels.*
fun main() = runBlocking {
//sampleStart
val numbers = produceNumbers() // produces integers from 1 and on
val squares = square(numbers) // squares integers
repeat(5) {
println(squares.receive()) // print first five
}
println("Done!") // we are done
coroutineContext.cancelChildren() // cancel children coroutines
//sampleEnd
}
fun CoroutineScope.produceNumbers() = produce<Int> {
var x = 1
while (true) send(x++) // infinite stream of integers starting from 1
}
fun CoroutineScope.square(numbers: ReceiveChannel<Int>): ReceiveChannel<Int> = produce {
for (x in numbers) send(x * x)
}
Простые числа с использованием цепочки
Давайте доведём цепочки до логического завершения с примером генерации простых чисел с использованием цепочки сопрограмм. Мы начинаем с бесконечной последовательности чисел.
fun CoroutineScope.numbersFrom(start: Int) = produce<Int> {
var x = start
while (true) send(x++) // infinite stream of integers from start
}
Следующий этап цепочки фильтрует входящий поток чисел, удаляя все числа, которые делятся на данное простое число:
fun CoroutineScope.filter(numbers: ReceiveChannel<Int>, prime: Int) = produce<Int> {
for (x in numbers) if (x % prime != 0) send(x)
}
Теперь мы создаём нашу цепочку, начав с потока чисел с 2, взяв простое число из текущего канала и запустив новый этап цепочки для каждого найденного простого числа:
numbersFrom(2) -> filter(2) -> filter(3) -> filter(5) -> filter(7) ...
Следующий пример выводит первые десять простых чисел, выполняя всю цепочку в контексте основного потока. Поскольку все сопрограммы запускаются в области действия основной сопрограммы runBlocking, нам не нужно хранить явный список всех запущенных сопрограмм. Мы используем функцию расширения cancelChildren для отмены всех дочерних сопрограмм после вывода первых десяти простых чисел.
import kotlinx.coroutines.*
import kotlinx.coroutines.channels.*
fun main() = runBlocking {
//sampleStart
var cur = numbersFrom(2)
repeat(10) {
val prime = cur.receive()
println(prime)
cur = filter(cur, prime)
}
coroutineContext.cancelChildren() // cancel all children to let main finish
//sampleEnd
}
fun CoroutineScope.numbersFrom(start: Int) = produce<Int> {
var x = start
while (true) send(x++) // infinite stream of integers from start
}
fun CoroutineScope.filter(numbers: ReceiveChannel<Int>, prime: Int) = produce<Int> {
for (x in numbers) if (x % prime != 0) send(x)
}
Вывод этого кода:
2 3 5 7 11 13 17 19 23 29
Обратите внимание, что вы можете построить ту же цепочку, используя iterator конструктор сопрограмм из стандартной библиотеки. Замените produce на iterator, send на yield, receive на next, ReceiveChannel на Iterator, и избавитесь от области видимости сопрограмм. Вам также не потребуется runBlocking. Однако преимущество цепочки, использующей каналы, как показано выше, заключается в том, что она может фактически использовать несколько ядер процессора, если вы запустите её в контексте Dispatchers.Default.
В любом случае, это крайне непрактичный способ поиска простых чисел. На практике цепочки обычно включают в себя другие приостанавливаемые вызовы (например, асинхронные вызовы к удалённым службам), и эти цепочки нельзя построить с помощью sequence/iterator, так как они не допускают произвольной приостановки, в отличие от produce, которая полностью асинхронна.
Раздача
Несколько сопрограмм могут получать данные из одного канала, распределяя работу между собой. Начнём с сопрограммы-производителя, которая периодически производит целые числа (по десять чисел в секунду):
fun CoroutineScope.produceNumbers() = produce<Int> {
var x = 1 // start from 1
while (true) {
send(x++) // produce next
delay(100) // wait 0.1s
}
}
Затем у нас могут быть несколько сопрограмм-обработчиков. В этом примере они просто выведут свой идентификатор и полученное число:
fun CoroutineScope.launchProcessor(id: Int, channel: ReceiveChannel<Int>) = launch {
for (msg in channel) {
println("Processor #$id received $msg")
}
}
Теперь запустим пять процессоров и позволим им работать почти секунду. Посмотрите, что произойдёт:
import kotlinx.coroutines.*
import kotlinx.coroutines.channels.*
fun main() = runBlocking<Unit> {
//sampleStart
val producer = produceNumbers()
repeat(5) { launchProcessor(it, producer) }
delay(950)
producer.cancel() // cancel producer coroutine and thus kill them all
//sampleEnd
}
fun CoroutineScope.produceNumbers() = produce<Int> {
var x = 1 // start from 1
while (true) {
send(x++) // produce next
delay(100) // wait 0.1s
}
}
fun CoroutineScope.launchProcessor(id: Int, channel: ReceiveChannel<Int>) = launch {
for (msg in channel) {
println("Processor #$id received $msg")
}
}
Вывод будет похож на следующий, хотя идентификаторы процессоров, которые получают каждое конкретное число, могут быть разными:
Processor #2 received 1 Processor #4 received 2 Processor #0 received 3 Processor #1 received 4 Processor #3 received 5 Processor #2 received 6 Processor #4 received 7 Processor #0 received 8 Processor #1 received 9 Processor #3 received 10
Обратите внимание, что отмена сопрограммы-производителя закрывает её канал, что в конечном итоге приводит к прекращению итерации по каналу, которую выполняют сопрограммы-обработчики.
Также обратите внимание на то, как мы явно итерируем по каналу с помощью цикла for для выполнения раздачи в коде launchProcessor. В отличие от consumeEach, этот шаблон цикла for идеально подходит для использования несколькими сопрограммами. Если одна из сопрограмм-обработчиков завершится неудачно, другие всё равно будут обрабатывать канал, в то время как сопрограмма-обработчик, написанная с помощью consumeEach, всегда потребляет (отменяет) подлежащий канал при нормальном или ненормальном завершении.
Слияние
Несколько сопрограмм могут отправлять данные в один и тот же канал. Например, у нас есть канал строк, и приостанавливаемая функция, которая многократно отправляет указанную строку в этот канал с указанной задержкой:
suspend fun sendString(channel: SendChannel<String>, s: String, time: Long) {
while (true) {
delay(time)
channel.send(s)
}
}
Теперь посмотрим, что произойдёт, если мы запустим несколько сопрограмм, отправляющих строки (в этом примере мы запускаем их в контексте основного потока как дочерние сопрограммы основной сопрограммы):
import kotlinx.coroutines.*
import kotlinx.coroutines.channels.*
fun main() = runBlocking {
//sampleStart
val channel = Channel<String>()
launch { sendString(channel, "foo", 200L) }
launch { sendString(channel, "BAR!", 500L) }
repeat(6) { // receive first six
println(channel.receive())
}
coroutineContext.cancelChildren() // cancel all children to let main finish
//sampleEnd
}
suspend fun sendString(channel: SendChannel<String>, s: String, time: Long) {
while (true) {
delay(time)
channel.send(s)
}
}
Вывод:
foo foo BAR! foo foo BAR!
Буферизованные каналы
Рассмотренные ранее каналы не имели буфера. Небуферизованные каналы передают элементы, когда отправитель и получатель встречаются (также известен как rendez-vous). Если сначала вызывается send, то оно приостанавливается до вызова receive, если сначала вызывается receive, то оно приостанавливается до вызова send.
Оба фабричные функции Channel() и билдер produce принимают необязательный capacity параметр для указания размера буфера. Буфер позволяет отправителям отправлять несколько элементов до приостановки, аналогично BlockingQueue с заданной ёмкостью, которая блокируется, когда буфер заполнен.
Посмотрите на поведение следующего кода:
import kotlinx.coroutines.*
import kotlinx.coroutines.channels.*
fun main() = runBlocking<Unit> {
//sampleStart
val channel = Channel<Int>(4) // create buffered channel
val sender = launch { // launch sender coroutine
repeat(10) {
println("Sending $it") // print before sending each element
channel.send(it) // will suspend when buffer is full
}
}
// don't receive anything... just wait....
delay(1000)
sender.cancel() // cancel sender coroutine
//sampleEnd
}
Он печатает "отправка" пять раз, используя буферизованный канал с ёмкостью четыре:
Sending 0 Sending 1 Sending 2 Sending 3 Sending 4
Первые четыре элемента добавляются в буфер, и отправитель приостанавливается при попытке отправить пятый.
Каналы справедливы
Операции отправки и получения в каналах являются справедливыми относительно порядка их вызова из нескольких сопроцессов. Они обрабатываются в порядке очереди «первым пришёл – первым обслужен», например, первый сопроцесс, вызывающий receive, получает элемент. В следующем примере две сопроцессы «ping» и «pong» получают объект «мяч» из общего канала «таблица».
import kotlinx.coroutines.*
import kotlinx.coroutines.channels.*
//sampleStart
data class Ball(var hits: Int)
fun main() = runBlocking {
val table = Channel<Ball>() // a shared table
launch { player("ping", table) }
launch { player("pong", table) }
table.send(Ball(0)) // serve the ball
delay(1000) // delay 1 second
coroutineContext.cancelChildren() // game over, cancel them
}
suspend fun player(name: String, table: Channel<Ball>) {
for (ball in table) { // receive the ball in a loop
ball.hits++
println("$name $ball")
delay(300) // wait a bit
table.send(ball) // send the ball back
}
}
//sampleEnd
Первой запускается сопроцесс «ping», поэтому он первым получает мяч. Даже если сопроцесс «ping» немедленно начинает снова получать мяч после возвращения его на стол, мяч получает сопроцесс «pong», потому что он уже ждал его:
ping Ball(hits=1) pong Ball(hits=2) ping Ball(hits=3) pong Ball(hits=4)
Обратите внимание, что иногда каналы могут производить выполнения, которые выглядят несправедливыми из-за природы используемого исполнителя. Подробнее см. эту проблему.
Каналы тикеров
Канал тикера — это специальный канал rendez-vous, который производит Unit каждый раз, когда заданная задержка проходит с момента последнего потребления из этого канала. Хотя он может показаться бесполезным самостоятельно, он является полезным строительным блоком для создания сложных каналов на основе времени produce и операторов, выполняющих оконную обработку и другие задачи обработки, зависящие от времени. Канал тикера можно использовать в select для выполнения действия «при тике».
Для создания такого канала используйте фабричный метод ticker. Чтобы указать, что больше элементов не требуется, используйте метод ReceiveChannel.cancel на нем.
Теперь давайте посмотрим, как это работает на практике:
import kotlinx.coroutines.*
import kotlinx.coroutines.channels.*
//sampleStart
fun main() = runBlocking<Unit> {
val tickerChannel = ticker(delayMillis = 100, initialDelayMillis = 0) // create ticker channel
var nextElement = withTimeoutOrNull(1) { tickerChannel.receive() }
println("Initial element is available immediately: $nextElement") // no initial delay
nextElement = withTimeoutOrNull(50) { tickerChannel.receive() } // all subsequent elements have 100ms delay
println("Next element is not ready in 50 ms: $nextElement")
nextElement = withTimeoutOrNull(60) { tickerChannel.receive() }
println("Next element is ready in 100 ms: $nextElement")
// Emulate large consumption delays
println("Consumer pauses for 150ms")
delay(150)
// Next element is available immediately
nextElement = withTimeoutOrNull(1) { tickerChannel.receive() }
println("Next element is available immediately after large consumer delay: $nextElement")
// Note that the pause between `receive` calls is taken into account and next element arrives faster
nextElement = withTimeoutOrNull(60) { tickerChannel.receive() }
println("Next element is ready in 50ms after consumer pause in 150ms: $nextElement")
tickerChannel.cancel() // indicate that no more elements are needed
}
//sampleEnd
Он выводит следующие строки:
Initial element is available immediately: kotlin.Unit Next element is not ready in 50 ms: null Next element is ready in 100 ms: kotlin.Unit Consumer pauses for 150ms Next element is available immediately after large consumer delay: kotlin.Unit Next element is ready in 50ms after consumer pause in 150ms: kotlin.Unit
Обратите внимание, что ticker учитывает возможные паузы потребителя и по умолчанию корректирует задержку следующего произведённого элемента в случае возникновения паузы, пытаясь поддерживать фиксированную скорость произведённых элементов.
Необязательно, параметр mode равный TickerMode.FIXED_DELAY может быть указан для поддержания фиксированной задержки между элементами.
© 2010–2023 JetBrains s.r.o. and Kotlin Programming Language contributors
Licensed under the Apache License, Version 2.0.
https://kotlinlang.org/docs/channels.html