1. Что такое Channel и зачем он нужен
Когда мы только начинаем писать корутинный код, рука тянется сделать так: «у меня есть producer-корутина, она добавляет элементы в MutableList, а consumer-корутина читает их оттуда». На маленьком примере это иногда даже “работает”, но ровно до момента, пока вы не встретите гонки данных, странные пропуски и непредсказуемые состояния.
Channel нужен как раз затем, чтобы вместо совместного владения изменяемыми структурами перейти к нормальному протоколу: один отправил сообщение, другой его получил.
Представим учебный консольный проект (пусть это будет утилита Console Monitor): одна корутина «собирает события» (например, строки логов или «псевдо-метрики»), а другая корутина эти события обрабатывает и печатает агрегированную информацию. Нам нужен безопасный «конвейер» между ними: не общий список, который все трогают руками, а нормальная очередь сообщений.
Схематично это выглядит так:
flowchart LR
P[Producer coroutine
генерирует события] --> C[(Channel⟨Event⟩)]
C --> R[Receiver coroutine
обрабатывает события]
Channel<T> как очередь сообщений с типом
В Kotlin корутинах Channel — это абстракция «очереди сообщений» между корутинами. Главное слово тут — сообщений. Мы не «делим память», а передаём значения.
Канал типизирован: Channel<Int> — передаёт Int, Channel<String> — строки, Channel<Event> — ваши собственные события.
Есть важная деталь, которая часто ломает ожидания новичков: канал — это не “рассылка всем подписчикам”. Одно отправленное значение будет доставлено ровно одному получателю (если получателей несколько, они “соревнуются” за сообщения). Это часть модели Channel.
Поэтому думайте о Channel как о «очереди задач/сообщений», а не как о «групповом чате», где одно сообщение должны увидеть все.
Мини-табличка для закрепления:
| Что мы хотим | Что лучше подходит | Интуиция |
|---|---|---|
| Передать значение “из корутины A в корутину B” | |
Как очередь: положили → забрали |
| Хранить набор данных “внутри программы” | коллекции ( , ) |
Как склад: лежит и ждёт |
| Раздать одно событие “всем слушателям” | не (это другой паттерн) |
Как рассылка/уведомления |
Про “рассылку всем” сегодня специально не углубляемся — наша цель научиться базовой модели: очередь send/receive.
Минимальный пример: send и receive
Самый маленький рабочий фрагмент, чтобы в голове щёлкнуло: канал — это штука, через которую реально можно «перекинуть значение» между корутинами.
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.launch
import kotlinx.coroutines.channels.Channel
fun main() = runBlocking {
val ch = Channel<Int>()
launch { ch.send(42) }
val x = ch.receive()
println("Got $x") // Got 42
ch.close()
}
Здесь важно не то, что мы напечатали 42, а то, что мы увидели новый стиль мышления: вместо «общей переменной» у нас появился канал связи.
И сразу маленькое наблюдение: send(...) и receive() работают только в корутинном контексте. Это не прихоть авторов библиотеки, а логика самого механизма: эти операции могут “ждать”, а значит должны быть suspend.
Почему send() и receive() — suspend
Если раньше вы воспринимали suspend как “внутри будет delay”, то здесь расширяем кругозор: suspend — это про операции, которые могут приостановить корутину без блокировки потока.
У канала оба конца потенциально «ожидающие»: отправитель может ждать, пока появится место, а получатель может ждать, пока появится значение.
Давайте это разложим в понятную механику. Канал можно представить как почтовый ящик между двумя людьми. Если ящик пуст, получатель заглянул — и ждёт, пока туда положат письмо. Если ящик переполнен (или вообще без буфера), отправитель хочет положить письмо — и ждёт, пока кто-то заберёт.
Таблица “кто когда ждёт”:
| Операция | Почему может приостановиться | Короткая человеческая фраза |
|---|---|---|
|
“Некуда положить” (нет места в буфере или буфер 0) | «Я подержу сообщение в руках, пока освободится место» |
|
“Нечего взять” (канал пуст) | «Я подожду, пока кто-то пришлёт сообщение» |
На этом месте часто появляется здравая мысль: «То есть канал сам замедляет отправителя, если получатель не успевает?» Да, и это очень полезное свойство. В мире потоков данных это часто называют “естественным тормозом”, чтобы producer не “залил” память миллионом сообщений.
Capacity: буфер канала и поведение send
Теперь добавим к “почтовому ящику” размер. В Channel это называется capacity (ёмкость).
Если capacity равна 0, то отправитель и получатель встречаются “лицом к лицу”: отправка завершится только когда есть получатель. Это похоже на передачу пакета из рук в руки.
Если capacity больше нуля, канал начинает работать как настоящая очередь: можно положить несколько элементов “вперёд”, не дожидаясь, пока их прямо сейчас заберут.
Пример с capacity = 1, где второе send может “подвиснуть”, если получатель не успел забрать первое:
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.launch
import kotlinx.coroutines.delay
import kotlinx.coroutines.channels.Channel
fun main() = runBlocking {
val ch = Channel<String>(capacity = 1)
val sender = launch {
ch.send("A")
ch.send("B") // может подождать, пока "A" не заберут
ch.close()
}
delay(50)
println(ch.receive()) // A
println(ch.receive()) // B
sender.join()
}
В этом примере задержка delay(50) делает ситуацию наглядной: producer пытается быстро отправить "A" и "B", но канал позволяет “протолкнуть” вперёд только один элемент. Второй send будет ждать, пока consumer освободит место.
Полезная аналогия: capacity — это размер “коробки у двери”. Если коробка на 1 посылку, курьер может оставить одну и уйти. Но если он пришёл с двумя, вторую придётся держать в руках, пока вы не заберёте первую.
close(): как сказать “сообщений больше не будет”
Любая очередь сообщений упирается в простой вопрос: “а когда заканчиваем?”. Если consumer делает receive() в цикле, но producer уже давно закончил — consumer может ждать вечно. Поэтому у канала есть важная часть протокола: close().
close() означает: “новых элементов больше не будет”. При этом те элементы, которые уже лежат в канале (в буфере), всё равно можно дочитать. Закрытие — это не “стереть и забыть”, это именно сигнал завершения потока сообщений.
Самый удобный и читаемый способ читать канал до закрытия — цикл for (x in ch), который сам корректно завершится, когда канал закрыт и элементы кончились.
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.launch
import kotlinx.coroutines.channels.Channel
fun main() = runBlocking {
val ch = Channel<Int>()
launch {
repeat(3) { i -> ch.send(i) }
ch.close()
}
for (x in ch) {
println("x=$x") // x=0, x=1, x=2
}
println("Done") // Done
}
Такой for работает именно потому, что канал — это поток значений “во времени”, и завершение у него выражается закрытием.
2. Практика: Console Monitor и протокол завершения
Канал событий
Соберём кусочки в маленькую “историю”, чтобы было похоже не на лабораторные 42, а на реальный код.
Пусть наш Console Monitor делает две вещи: генерирует события (как будто это входящие строки логов), и печатает их, как будто это обработчик. Мы пока не строим “фабрику воркеров” и не делаем сложные пайплайны — нам важна базовая связка Channel + send/receive + close.
Сначала договоримся о типе события. Чтобы не усложнять лекцию, возьмём просто строку:
import kotlinx.coroutines.channels.Channel
typealias LogEvent = String
fun createLogChannel(): Channel<LogEvent> = Channel(capacity = 2)
Producer: генерация событий и close
Теперь сделаем “генератор событий”: он в цикле отправляет несколько строк и закрывает канал.
import kotlinx.coroutines.delay
import kotlinx.coroutines.channels.Channel
suspend fun produceLogs(ch: Channel<String>) {
repeat(5) { i ->
ch.send("event#$i")
delay(30)
}
ch.close()
}
Consumer: чтение до закрытия
А теперь “обработчик событий”: он читает всё, пока канал не закрыт.
import kotlinx.coroutines.channels.Channel
suspend fun consumeLogs(ch: Channel<String>) {
for (e in ch) {
println("handle: $e") // handle: event#0 ...
}
}
Склейка в main
Осталось соединить в main:
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.launch
import kotlinx.coroutines.channels.Channel
fun main() = runBlocking {
val ch = Channel<String>(capacity = 2)
val producer = launch { produceLogs(ch) }
val consumer = launch { consumeLogs(ch) }
producer.join()
consumer.join()
}
Да, это похоже на “игрушку”. Но это правильная игрушка: в ней уже есть протокол завершения (close), и нет общей мутабельной коллекции, которую две корутины терзают одновременно.
close в finally: чтобы consumer не зависал при ошибках
На практике producer может завершиться не только “по плану”. Может упасть исключение, может сработать отмена, может вы выкинули return раньше времени. Если в таких сценариях канал не закрыть, consumer может зависнуть на ожидании новых сообщений, потому что с его точки зрения “жизнь ещё не закончилась”.
Поэтому полезная дисциплина: если именно ваш код отвечает за завершение канала, закрывайте канал в finally. Это ровно та же идея, что и с любыми ресурсами: “прибраться гарантированно”.
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.launch
import kotlinx.coroutines.channels.Channel
fun main() = runBlocking {
val ch = Channel<Int>()
val producer = launch {
try {
repeat(3) { i -> ch.send(i) }
} finally {
ch.close()
}
}
for (x in ch) println("x=$x") // x=0, x=1, x=2
producer.join()
}
Здесь даже если внутри try случится что-то неприятное, finally всё равно выполнится, канал будет закрыт, и цикл for (x in ch) корректно завершится.
3. Типичные ошибки при работе с Channel
Ошибка №1: забыли close() — и программа “висит”, как будто задумалась о смысле жизни.
Это самый частый сценарий: consumer читает из канала в цикле, а producer закончил работу и ушёл, не закрыв канал. Потребитель честно ждёт следующий элемент и будет ждать до конца времён (или пока вы не нажмёте Stop в IDE). Лечится не магией, а протоколом: канал закрывает тот, кто “владеет производством” сообщений, и часто это делается в finally.
Ошибка №2: попытка вызвать send/receive из обычной функции без suspend и без корутинного контекста.
send и receive могут приостанавливать корутину, поэтому Kotlin требует, чтобы вы находились в runBlocking или внутри launch/async, или хотя бы в suspend-функции, вызываемой из корутины. Если попытаться “просто вызвать”, компилятор будет ругаться, и это тот редкий случай, когда компилятор не зануда, а ваш телохранитель.
Ошибка №3: ожидание, что один Channel разошлёт одно сообщение всем потребителям.
Интуиция “канал = трансляция” встречается часто, но модель другая: каждое значение из Channel получает только один получатель. Если у вас два consumer-а, они будут делить сообщения между собой, а не дублировать. Это не баг — это главный смысл очереди задач.
Ошибка №4: путаница между capacity и “количеством получателей”.
Ёмкость канала влияет только на буферизацию: сколько элементов можно временно “положить внутрь” без ожидания отправителя. Она никак не говорит, сколько корутин “увидят” сообщение. Сообщение всё равно будет доставлено ровно одному получателю, просто отправка может происходить либо “рука-в-руку”, либо с небольшим запасом буфера.
Ошибка №5: закрывают канал “где попало”, а потом удивляются исключениям при send.
Если вы закрыли канал, а кто-то потом пытается отправить ещё одно сообщение, это уже нарушение протокола. В результате вы получите падение (или, как минимум, очень неприятное поведение). Выбирайте одно место, где канал закрывается, и делайте это место логически ответственным за завершение — иначе код превращается в детектив, где убийца “кто-то из нас”.
Ошибка №6: чтение канала через бесконечный while (true) { receive() } без нормального условия завершения.
Технически так написать можно, но это почти всегда усложняет жизнь: вам придётся отдельно думать, как вы выйдете из цикла, что будет на закрытом канале и где ловить исключение. Цикл for (x in ch) читается проще и напрямую отражает смысл: “читай всё, пока не кончится”.
ПЕРЕЙДИТЕ В ПОЛНУЮ ВЕРСИЮ