Каналы
Отложенные значения предоставляют удобный способ передать одно значение между корутинами. Каналы позволяют передавать поток значений.
Основы каналов
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.*
//sampleStart
fun CoroutineScope.produceSquares(): ReceiveChannel<Int> = produce {
for (x in 1..5) send(x * x)
}
fun main() = runBlocking {
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
Обратите внимание: отмена корутины-производителя закрывает её канал, что в конечном итоге приводит к завершению перебора канала корутинами-обработчиками.
Также обратите внимание, что для распределения нагрузки в коде launchProcessor мы явно перебираем канал с помощью цикла for. В отличие от 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
Первые четыре элемента добавляются в буфер, а при попытке отправить пятый отправитель приостанавливается.
Каналы справедливы
Операции отправки и получения данных из каналов справедливы по отношению к порядку их вызова из нескольких корутин. Они обслуживаются в порядке «первым пришёл — первым обслужен»: например, элемент получает первая корутина, вызвавшая 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 каждый раз по истечении заданной задержки с момента последнего получения данных из канала. Хотя сам по себе он может показаться бесполезным, это удобный строительный блок для создания сложных временных конвейеров produce и операторов, выполняющих оконную и другую зависящую от времени обработку. Тикер-канал можно использовать в select для выполнения действия «по тику».
Чтобы создать такой канал, используйте фабричный метод ticker. Чтобы указать, что больше элементы не нужны, вызовите для него метод ReceiveChannel.cancel.
Посмотрим, как это работает на практике:
import kotlinx.coroutines.*
import kotlinx.coroutines.channels.*
//sampleStart
fun main() = runBlocking<Unit> {
val tickerChannel = ticker(delayMillis = 200, initialDelayMillis = 0) // create a ticker channel
var nextElement = withTimeoutOrNull(1) { tickerChannel.receive() }
println("Initial element is available immediately: $nextElement") // no initial delay
nextElement = withTimeoutOrNull(100) { tickerChannel.receive() } // all subsequent elements have 200ms delay
println("Next element is not ready in 100 ms: $nextElement")
nextElement = withTimeoutOrNull(120) { tickerChannel.receive() }
println("Next element is ready in 200 ms: $nextElement")
// Emulate large consumption delays
println("Consumer pauses for 300ms")
delay(300)
// 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(120) { tickerChannel.receive() }
println("Next element is ready in 100ms after consumer pause in 300ms: $nextElement")
tickerChannel.cancel() // indicate that no more elements are needed
}
//sampleEnd
Выводятся следующие строки:
Initial element is available immediately: kotlin.Unit Next element is not ready in 100 ms: null Next element is ready in 200 ms: kotlin.Unit Consumer pauses for 300ms Next element is available immediately after large consumer delay: kotlin.Unit Next element is ready in 100ms after consumer pause in 300ms: kotlin.Unit
Обратите внимание, что ticker учитывает возможные паузы потребителя и по умолчанию корректирует задержку перед следующим элементом, если пауза возникла, стараясь поддерживать постоянную частоту выдачи элементов.
При необходимости можно указать параметр mode со значением TickerMode.FIXED_DELAY, чтобы поддерживать фиксированную задержку между элементами.
© 2010–2026 JetBrains s.r.o. and Kotlin Programming Language contributors
Licensed under the Apache License, Version 2.0.
https://kotlinlang.org/docs/channels.html