Каналы
Отложенные значения предоставляют удобный способ передачи одного значения между сопрограммами. Каналы обеспечивают способ передачи потока значений.
Основы каналов
Канал 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!
Буферизованные каналы
Продемонстрированные ранее каналы не имели буфера. Небуферизованные каналы передают элементы, когда отправитель и получатель встречаются (т.е. согласование). Если сначала вызывается отправка, она приостанавливается до вызова приема; если сначала вызывается прием, он приостанавливается до вызова отправки.
Оба фабричные функции 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" пять раз, используя буферизованный канал с ёмкостью четыре:
Sending 0 Sending 1 Sending 2 Sending 3 Sending 4
Первые четыре элемента добавляются в буфер, а отправитель приостанавливается при попытке отправить пятый.
Каналы справедливы
Операции отправки и получения в каналах являются справедливыми относительно порядка их вызова из нескольких сопроцедур. Они обрабатываются в порядке очереди «первый пришёл — первый обслужен», например, первая сопроцедура, вызвавшая receive, получает элемент. В следующем примере две сопроцедуры «ping» и «pong» получают объект «ball» из общего канала «table».
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)
Обратите внимание, что иногда каналы могут генерировать операции, которые выглядят несправедливыми из-за характера используемого исполнителя. Подробнее об этом можно узнать в этом вопросе.
Каналы таймера
Канал таймера — это специальный канал согласования, который производит 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–2022 JetBrains s.r.o. and Kotlin Programming Language contributors
Licensed under the Apache License, Version 2.0.
https://kotlinlang.org/docs/channels.html