Spec-Zone.ru › Kotlin 1.4

Содержание

  • Каналы
    • Основы каналов
    • Закрытие и итерация по каналам
    • Создание производителей каналов
    • Потоки
    • Простые числа с потоком
    • Fan-out
    • Fan-in
    • Буферизованные каналы
    • Каналы честные
    • Каналы-таймеры

Каналы

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

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

Канал 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, который полностью асинхронен.

Fan-out

Несколько корутин могут получать из одного канала, распределяя работу между собой. Давайте начнём с корутины-производителя, которая периодически генерирует целые числа (по десять чисел в секунду):

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 циклом, чтобы выполнить fan-out в коде launchProcessor. В отличие от consumeEach, этот for паттерн цикла совершенно безопасен для использования из нескольких корутин. Если одна из корутин-процессоров завершится ошибкой, то другие всё ещё будут обрабатывать канал, в то время как процессор, написанный через consumeEach, всегда потребляет (отменяет) подчинённый канал при нормальном или аномальном завершении.

Fan-in

Несколько корутин могут отправлять в один и тот же канал. Например, у нас есть канал строк, и приостанавливаемая функция, которая многократно отправляет указанную строку в этот канал с указанной задержкой:

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!

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

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

Функция-фабрика 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)

Обратите внимание, что иногда каналы могут выдавать результаты, которые кажутся несправедливыми из-за характера используемого исполнителя. Подробности см. в этой проблеме.

Каналы таймера

Канал таймера — это специальный канал rendez-vous, который производит Unit каждый раз, когда заданная задержка проходит с момента последнего использования этого канала. Хотя сам по себе он может показаться бесполезным, он является полезным строительным блоком для создания сложных каналов, основанных на времени, и операторов produce, которые выполняют оконную обработку и другие зависящие от времени операции. Канал таймера может быть использован в 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–2020 JetBrains s.r.o. and Kotlin Programming Language contributors
Licensed under the Apache License, Version 2.0.
https://kotlinlang.org/docs/reference/coroutines/channels.html

Spec-Zone.ru

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