Каналы
Отложенные значения предоставляют удобный способ передачи одного значения между сопроцедурами. Каналы предоставляют способ передачи потока значений.
Основы каналов
Канал 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!
Буферизованные каналы
Рассмотренные ранее каналы не имели буфера. Небуферизованные каналы передают элементы, когда отправитель и получатель встречаются (также известен как rendezvous). Если сначала вызывается 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" пять раз, используя буферизованный канал с емкостью четыре:
Sending 0 Sending 1 Sending 2 Sending 3 Sending 4
Первые четыре элемента добавляются в буфер, а отправитель приостанавливается при попытке отправить пятый.
Каналы справедливы
Операции отправки и получения в каналах являются справедливыми по отношению к порядку их вызова из нескольких сопроцессов. Они обслуживаются в порядке очереди FIFO, например, первый сопроцесс, вызывающий 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)
Обратите внимание, что иногда каналы могут производить вычисления, которые кажутся несправедливыми из-за характера используемого исполнителя. Подробнее об этом см. в этой проблеме.
Каналы-таймеры
Канал-таймер — это специальный канал rendezvous, который производит Unit каждый раз, когда заданная задержка проходит с момента последнего потребления из этого канала. Хотя он может показаться бесполезным сам по себе, он является полезным строительным блоком для создания сложных каналов и операторов на основе времени, которые выполняют оконную обработку и другие операции, зависящие от времени. Канал-таймер может быть использован в select для выполнения действия «при срабатывании таймера».
Для создания такого канала используйте фабричный метод ticker. Для указания, что больше элементов не требуется, используйте метод ReceiveChannel.cancel на нём.
Теперь давайте посмотрим, как он работает на практике:
import kotlinx.coroutines.*
import kotlinx.coroutines.channels.*
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
}
Он выводит следующие строки:
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