Spec-Zone.ru › Kotlin 1.7

Каналы

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

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

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

Буферизованные каналы

Продемонстрированные ранее каналы не имели буфера. Небуферизованные каналы передают элементы, когда отправитель и получатель встречаются (т.е. согласование). Если сначала вызывается отправка, она приостанавливается до вызова приема; если сначала вызывается прием, он приостанавливается до вызова отправки.

Оба фабричные функции 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, чтобы поддерживать постоянную задержку между элементами.

Последнее изменение: 27 июня 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