Spec-Zone.ru › Kotlin 1.6

Каналы

Отложенные значения предоставляют удобный способ передачи одного значения между сопроцедурами. Каналы предоставляют способ передачи потока значений.

Основы каналов

Канал 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)
}

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

Все функции, создающие сопроцедуры, определены как расширения для CoroutineScope, так что мы можем полагаться на структурированную конкурентность, чтобы убедиться, что в нашем приложении нет незавершенных глобальных сопроцедур.

Простое число с потоком данных

Давайте доведем потоки данных до крайности на примере генерации простых чисел с помощью потока сопроцедур. Мы начнем с бесконечной последовательности чисел.

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 для поддержания постоянной задержки между элементами.

Последнее изменение: 04 апреля 2022
Асинхронный поток Обработка исключений сопроцессов

© 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

Spec-Zone.ru

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