Тема 1. Как выглядит Kotlin Coroutine без макияжа
Тема 2. Kotlin suspend функции
Тема 3. Kotlin Coroutine диспетчеры и потоки: где выполняются корутины?
Тема 4. Structured Concurrency и отмена корутин
Тема 5. Обработка исключений в Kotlin Coroutines

Пытаюсь лучше понять работу Kotlin Coroutines. Беру небольшую тему про Kotlin Coroutines и пытаюсь разобраться и написать шпаргалку (максимально кратко и лаконично). Сегодня поговорим про Каналы (Channel) и паттерн Actor.

В предыдущих частях я писал, что обычные suspend‑функции возвращают результат только один раз. Но что делать, если у нас есть непрерывный поток данных, который асинхронно генерируется в одной корутине, а обрабатываться должен в другой (или даже в нескольких)? Нам нужна потокобезопасная «труба». В мире корутин этой трубой является Channel.

1. Что такое Channel под капотом?

Концептуально Channel — это потокобезопасная очередь (Queue), адаптированная специально для корутин.

Главное отличие канала от классических блокирующих очередей заключается в том, что его методы не блокируют потоки операционной системы. Вместо этого они используют механизм приостановки. Как мы помним, suspend‑функции позволяют приостанавливать выполнение кода без блокировки потока.

Интерфейс канала логически разделен на две части для соблюдения инкапсуляции:

  1. SendChannel — интерфейс для отправителя (имеет suspend‑функцию send()).

  2. ReceiveChannel — интерфейс для получателя (имеет suspend‑функцию receive()).

Когда корутина‑отправитель вызывает send(), а получатель еще не готов (например, буфер заполнен), отправитель просто приостанавливается (suspend), освобождая текущий поток для других задач. Состояние корутины сохраняется в Continuation, пока получатель не заберет данные. То же самое работает и в обратную сторону для receive().

2. Виды каналов (Capacity)

При создании канала мы можем управлять размером его внутреннего буфера: Channel<T>(capacity). От этого параметра кардинально меняется поведение системы:

  • Rendezvous (по умолчанию, capacity = 0): Буфера нет вообще. Отправитель и получатель должны «встретиться». Функция send() приостановит корутину и будет ждать, пока на другой стороне не вызовут receive().

  • Buffered (capacity = N): Канал имеет массив заданного размера N. send() отрабатывает мгновенно и не приостанавливает корутину, пока буфер не заполнится.

  • Conflated (capacity = 1): Хранит только последнее отправленное значение. Если получатель не успевает читать, старые данные просто перезаписываются новыми. send() здесь никогда не приостанавливает корутину.

  • Unlimited: Безлимитный буфер. send() никогда не уходит в suspend, но если отправитель работает быстрее получателя, приложение неминуемо упадет с OutOfMemoryError.

3. Backpressure

Одна из главных проблем в асинхронном программировании — разница скоростей. Что будет, если источник данных генерирует их с огромной скоростью, а база данных не успевает их сохранять?

Здесь каналы показывают свою магию «из коробки» — механизм Backpressure. Если мы используем Buffered канал, то при заполнении буфера очередная попытка вызвать send() просто приостановит корутину‑продюсера (отправитель). Она «уснёт», перестав потреблять процессорное время и память, пока консьюмер (получатель) не разгребет очередь и не освободит место в буфере.

4. Практический пример:

Представим ситуацию: мы разрабатываем мобильное приложение, которое общается с внешним измерительным оборудованием. Каждую миллисекунду по BLE GATT или Wi‑Fi (UDP) нам прилетает непрерывный поток байтов, которые нужно склеивать, парсить и складывать в локальную БД.

Если использовать обычные коллбеки, мы можем забить память или потерять пакеты из‑за состояния гонки. Как это решают каналы:

// Создаем буферизированный канал для сглаживания пиковых нагрузок
val packetChannel = Channel<ByteArray>(capacity = 64)

// Корутина 1: Продюсер (слушает железо)
launch(Dispatchers.IO) {
    bleManager.observeData { rawBytes ->
        // Если буфер забьется, send() приостановит получение новых данных, 
        // давая консьюмеру время на обработку
        packetChannel.send(rawBytes) 
    }
}

// Корутина 2: Консьюмер (обрабатывает данные)
launch(Dispatchers.Default) {
    // for-цикл по каналу под капотом вызывает receive()
    for (packet in packetChannel) { 
        val domainModel = parsePacket(packet)
        database.save(domainModel) // тяжелая suspend-операция
    }
}

Здесь каналы показывают себя достаточно эффективно, так как встроенный буфер позволяет сглаживать пиковые нагрузки

5. Actor‑модель: Управление состоянием без Mutex

Каналы умеют не только гонять данные, но и решать классическую проблему многопоточности — состояние гонки (Race Condition), когда несколько корутин пытаются изменить одну переменную.

Обычно для защиты разделяемого состояния используют мьютексы (Mutex). Но есть другой архитектурный подход — Actor‑модель.

Актор — это корутина, которая инкапсулирует внутри себя какое‑то состояние (переменные) и никому не дает к нему прямого доступа. Взаимодействие с этим состоянием происходит исключительно через Channel. Поскольку канал обрабатывает входящие сообщения строго последовательно (одно за другим в цикле for), состояние внутри актора меняется потокобезопасно, без необходимости использовать блокировки или Mutex.

Представим работу с Bluetooth‑устройством (BLE GATT). Разные экраны и фоновые процессы могут параллельно пытаться отправить байты в сокет. Если сделать это без синхронизации, байты перемешаются, и протокол сломается.

Опишем команды для нашего актора:

sealed class BleMsg {
    class SendCommand(val bytes: ByteArray) : BleMsg()
    class GetStatus(val response: CompletableDeferred<Boolean>) : BleMsg()
    object Disconnect : BleMsg()
}

В стандартной библиотеке корутин есть билдер actor, который лаконично реализует этот паттерн:

fun CoroutineScope.bleDeviceActor(): SendChannel<BleMsg> 
  = actor<BleMsg>(Dispatchers.IO) {
    // Внутреннее состояние, защищенное актором. 
    // Доступ к нему имеет только эта корутина.
    var isConnected = true 
    
    // Актор последовательно вытаскивает сообщения из канала
    for (msg in channel) { 
        when (msg) {
            is BleMsg.SendCommand -> {
                // Пишем в сокет строго по очереди, никаких гонок
                if (isConnected) writeToGattSocket(msg.bytes) 
            }
            is BleMsg.GetStatus -> msg.response.complete(isConnected)
            is BleMsg.Disconnect -> {
                isConnected = false
                closeConnection()
            }
        }
    }
}

Примечание: функция actor сейчас находится в статусе Obsolete в стандартной библиотеке корутин, но сам архитектурный паттерн остается актуальным, и его легко реализовать вручную через обычный Channel и корутину:

fun CoroutineScope.manualBleDeviceActor(): SendChannel<BleMsg> {
    val channel = Channel<BleMsg>() // Создаем канал для приема сообщений
    
    // Запускаем корутину, которая будет эксклюзивно владеть состоянием
    launch(Dispatchers.IO) {
        var isConnected = true 
        
        for (msg in channel) { 
            when (msg) {
                is BleMsg.SendCommand -> {
                    if (isConnected) writeToGattSocket(msg.bytes)
                }
                
                is BleMsg.GetStatus -> {
                  msg.response.complete(isConnected)
                }
                
                is BleMsg.Disconnect -> {
                    isConnected = false
                    closeConnection()
                }
            }
        }
    }
    
    return channel
}

Теперь из любого места в приложении мы можем безопасно вызывать myBleActor.send(BleMsg.SendCommand(byteArray)). Мы получаем именно SendChannel и на нем не получится вызвать receive() и нарушить инкапсуляцию очереди. Канал выстроит все параллельные вызовы в аккуратную очередь, и актор обработает их строго последовательно, сохранив целостность протокола.

6. Отмена и закрытие каналов

Важно разделять два понятия: закрытие канала и отмену корутины. В статье про Structured Concurrency я писал, что отмена корутин работает через выброс исключения CancellationException.

Когда цикл for завершает работу? Смотря на примеры кода с циклом for (item in channel), может сложиться обманчивое впечатление, что как только элементы в буфере закончатся, цикл прервется и код пойдет дальше. Это не так. Канал по своей природе — это бесконечная труба. Если буфер пуст, цикл for (под капотом вызывающий receive()) не заканчивается. Он просто приостанавливает корутину‑получателя (переводит её в состояние suspend). Корутина будет «спать» и бесконечно ждать, пока на другом конце не появится новый элемент.

Что значит закрытие канала? Закрытие канала — это явный сигнал: «Новых данных больше никогда не будет». Вызывать channel.close() имеет смысл тогда, когда ваш поток данных имеет логический конец. Например, вы построчно вычитали весь файл или выкачали все страницы из пагинированного сетевого ответа. Инициатором закрытия обычно выступает отправитель.

Что происходит после вызова channel.close()?

  1. Для отправителя: Канал мгновенно «запечатывается» на вход. Любая попытка вызвать send() после закрытия приведет к крашу с ошибкой ClosedSendChannelException.

  2. Для получателя: Все данные, которые к моменту закрытия уже успели попасть в буфер канала, никуда не исчезают. Получатель продолжит спокойно их вычитывать.

  3. Завершение цикла: И вот только теперь, когда выполняются два условия одновременно (канал закрыт И его буфер полностью опустел), стандартный цикл for (item in channel) штатно прервется, и корутина‑консьюмер перейдет к выполнению следующего за циклом кода.

Если же получатель читает данные не через цикл for, а вручную вызывая receive(), то попытка прочитать элемент из закрытого и пустого канала выбросит ClosedReceiveChannelException.

Нужно ли закрывать каналы всегда? Нет, это не обязательно. Если ваш канал используется для бесконечных событий (например, поток координат GPS, нажатия кнопок пользователем или постоянное BLE‑соединение из примера выше), вам не нужно ломать голову над вызовом close(). Когда жизненный цикл экрана или фичи завершится, вы просто отменяете родительский CoroutineScope. Этот механизм каскадно отменит все работающие корутины (и отправителя, и получателя) через штатный CancellationException, а сам объект канала без проблем будет собран сборщиком мусора.

Итоги

  • Channel — это очередь, где вместо блокировки потоков используется suspend.

  • Capacity (размер буфера) определяет, когда отправитель будет приостановлен.

  • Механизм приостановки дает нам бесплатный Backpressure, не позволяя быстрому отправителю уронить приложение по памяти.

  • Actor‑модель — это паттерн, где состояние прячется внутри корутины, а доступ к нему идет через последовательную обработку сообщений из канала, что позволяет избежать гонок данных без Mutex.