JavaRush /Курсы /Kotlin SELF /Producer–Consumer на Channe...

Producer–Consumer на Channel

Kotlin SELF
54 уровень , 1 лекция
Открыта

1. Знакомимся с шаблоном Producer–Consumer

Когда вы только учитесь программировать, хочется решать всё просто: «а давайте положим задачи в MutableList, а воркеры пусть из неё берут». Логика понятная: сверху задача, снизу обработка, и никаких этих ваших «каналов».

Но как только появляется параллельность, MutableList превращается в клубок проблем. Кто-то добавил элемент, кто-то удалил, кто-то начал перебирать, кто-то «чуть-чуть» поменял счётчик — и вы получаете ошибки, которые возникают раз в 30 запусков (то есть идеально, чтобы испортить себе день).

Producer–Consumer полезен тем, что вы разделяете роли и договариваетесь, кто и как передаёт данные. Channel в этой картине выступает как «общая очередь», но без необходимости делить один и тот же изменяемый список руками. И это очень в стиле Kotlin: лучше сделать понятный контракт, чем надеяться на удачу.

2. Термины и роли

Слова Producer и Consumer звучат так, будто мы продаём подписку на кофе: «производитель удовольствия, потребитель радости». На практике всё прозаичнее: producer создаёт значения, consumer их обрабатывает.

Кстати, терминология «producer/consumer» встречается и в других местах Kotlin‑мира. Например, в дженериках есть удобная мнемоника «Producer — читает (даёт наружу), Consumer — пишет (принимает внутрь)». Это объясняется в документации про variance и PECS‑идею (Producer‑Extends, Consumer‑Super). Сегодня у нас «producer/consumer» не про типы, а про поток задач, но интуиция совпадает: producer «производит» задания, consumers «потребляют» их и делают работу.

Для нашей схемы обычно нужны четыре роли:

Роль Что делает Типичная корутина
Producer генерирует задачи и отправляет в tasks
launch { ... tasks.send(...) ... }
Worker (Consumer) читает задачи из tasks, выполняет обработку, отправляет в results
launch { for (t in tasks) results.send(...) }
Collector читает results и выводит/сохраняет
launch/for (r in results) или просто цикл в runBlocking
Coordinator (иногда отдельный) следит, чтобы каналы закрывались вовремя
launch { join(); close() }

Важно: воркеров может быть несколько, producer обычно один (но бывает и несколько), collector чаще всего один (но тоже бывает иначе).

3. Мини-проект TaskPipeline

Чтобы примеры не были «в вакууме», будем развивать маленькое консольное приложение TaskPipeline. Оно делает простую вещь: producer создаёт задачи (например, числа), несколько воркеров считают «результат обработки» (например, возводят в квадрат с задержкой), и мы печатаем итоги.

Почему задержка? Потому что без неё параллельность выглядит как магия: всё успевает мгновенно, и кажется, что launch ничего не делает. Поэтому будем вставлять delay(...), чтобы увидеть, как задачи реально распределяются.

Мини-модели данных

Мини‑модели данных (их удобно держать в одном файле рядом с main):

data class Task(val id: Int, val value: Int)
data class TaskResult(val taskId: Int, val workerId: Int, val output: Int)

Один канал задач и несколько воркеров

Ключевая идея Channel, которую нужно почувствовать руками: одно сообщение из канала забирает ровно один получатель. Это не рассылка «всем», это очередь «кому досталось — того и тапки».

Поэтому если мы сделаем 3 воркера, которые читают один и тот же Channel<Task>, задачи автоматически распределятся между ними: кто успел receive, тот и взял следующую работу.

Вот минимальный пример распределения. Обратите внимание: здесь пока нет канала результатов — мы просто печатаем, кто что обработал.

import kotlinx.coroutines.delay
import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.channels.Channel

fun main() = runBlocking {
    val tasks = Channel<Task>()

    repeat(2) { workerId ->
        launch {
            for (t in tasks) {
                delay(50) // имитируем работу
                println("worker=$workerId processed task=${t.id}") // например: worker=1 processed task=3
            }
        }
    }

    repeat(5) { id ->
        tasks.send(Task(id = id, value = id * 10))
    }
    tasks.close()
}

Здесь producer — это наш runBlocking (главная корутина), а два launch — consumers. Мы закрываем tasks, чтобы цикл for (t in tasks) завершился у каждого воркера, когда задачи закончатся.

Если не закрыть канал, воркеры будут ждать бесконечно. Это похоже на ситуацию, когда вы позвали друзей на вечеринку, но забыли сказать, что вечеринка вообще-то уже закончилась: они будут стоять у холодильника и надеяться.

4. Канал результатов и протокол завершения

Добавляем results

Печатать внутри воркера удобно для демонстрации, но в реальной жизни воркер обычно не «сам решает, куда выводить». Он просто делает работу и отдаёт результат дальше.

Поэтому вводим второй канал:

  • tasks: Channel<Task> — очередь заданий
  • results: Channel<TaskResult> — очередь результатов

И теперь воркер отправляет результат в results.

import kotlinx.coroutines.delay
import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.channels.Channel

fun main() = runBlocking {
    val tasks = Channel<Task>()
    val results = Channel<TaskResult>()

    repeat(2) { workerId ->
        launch {
            for (t in tasks) {
                delay(50)
                val out = t.value * t.value
                results.send(TaskResult(t.id, workerId, out))
            }
        }
    }

    launch {
        repeat(5) { id -> tasks.send(Task(id, id + 1)) }
        tasks.close()
    }

    // Пока просто читаем 5 результатов вручную (а чуть позже — красиво)
    repeat(5) {
        val r = results.receive()
        println("task=${r.taskId} by worker=${r.workerId} => ${r.output}")
        // например: task=0 by worker=1 => 4
    }
}

И тут возникает первая серьёзная проблема: кто и когда должен закрыть results?

Почему закрытие каналов — часть контракта

В системе Producer–Consumer закрытие каналов — это не «последняя строка программы», а часть контракта.

Логика такая:

  1. Producer перестал производить задачи → закрываем tasks, чтобы воркеры смогли завершиться.
  2. Воркеры завершились → они больше не будут отправлять результаты → теперь можно закрыть results, чтобы сборщик результатов тоже завершился.

Схема (очень упрощённо) выглядит так:

flowchart LR
    P[Producer] -->|send Task| T[(tasks Channel)]
    T -->|receive Task| W1[Worker 1]
    T -->|receive Task| W2[Worker 2]
    W1 -->|send Result| R[(results Channel)]
    W2 -->|send Result| R
    R -->|receive Result| C[Collector]

    P -->|close tasks| T
    W1 -->|finish| X1(( ))
    W2 -->|finish| X2(( ))
    X1 -->|after join workers| R
    X2 -->|close results| R

Главная тонкость: tasks.close() не означает «результаты тоже можно закрывать». Воркеры могут ещё обрабатывать уже полученные задачи и пытаться отправить результат. Если вы закрыли results слишком рано — воркер упадёт при send.

Поэтому появляется роль координатора: корутина (или просто кусок кода), который дождётся завершения тех, кто отправляет в канал, и только потом его закроет.

Каноническая сборка пайплайна

Соберём «канонический» (для нашего уровня) вариант, где:

  • producer создаёт задачи и закрывает tasks в finally;
  • воркеры читают tasks до закрытия и отправляют результаты;
  • отдельная корутина ждёт воркеров (join) и закрывает results;
  • главный поток читает results циклом for (r in results) до закрытия.
import kotlinx.coroutines.delay
import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.joinAll
import kotlinx.coroutines.channels.Channel

fun main() = runBlocking {
    val tasks = Channel<Task>()
    val results = Channel<TaskResult>()

    val workers = List(2) { workerId ->
        launch {
            for (t in tasks) {
                delay(50)
                results.send(TaskResult(t.id, workerId, t.value * t.value))
            }
        }
    }

    val producer = launch {
        try {
            repeat(5) { id -> tasks.send(Task(id, id + 1)) }
        } finally {
            tasks.close()
        }
    }

    val closer = launch {
        producer.join()
        workers.joinAll()
        results.close()
    }

    for (r in results) {
        println("task=${r.taskId} by worker=${r.workerId} => ${r.output}")
    }

    closer.join()
}

Обратите внимание на структуру завершения.

Мы не «угадываем», когда пора закрыть results. Мы ждём: сначала producer, потом всех workers. Только после этого закрываем results, и цикл for (r in results) сам закончится. Это и есть аккуратный протокол.

Правило «владелец закрытия»

Есть практическое правило, которое стоит запомнить как заклинание против зависаний:

Закрывает канал тот, кто отвечает за то, что в него больше никто не будет отправлять.

Звучит как занудство, но на деле это спасает. Если у вас несколько воркеров отправляют в results, то закрыть results должен не «любой из них по желанию», а тот, кто может гарантировать, что все закончили отправку. Обычно это координатор (как в примере выше), потому что именно он знает список Job воркеров и умеет дождаться join().

Если попытаться закрыть results внутри одного воркера, он закроет канал, пока другие воркеры ещё работают. Получится конфликт: один уже «закрыл дверь», другие всё ещё пытаются заносить коробки.

5. Порядок результатов и ёмкость канала

Почему порядок результатов не гарантирован

Вопрос, который часто возникает у новичков: «А результаты придут в том же порядке, что задачи?»

И честный ответ: не обязаны. Как только у вас несколько воркеров, «кто быстрее обработал — тот раньше прислал». Поэтому collector получает результаты в порядке готовности, а не в порядке task.id.

Иногда это идеально (например, вы показываете пользователю прогресс). Иногда вам нужно восстановить порядок. Для восстановления порядка обычно используют буферизацию результатов (например, MutableMap<taskId, result>) или сортировку после сбора — но это уже отдельная тема проектирования, и в рамках сегодняшней лекции нам важно хотя бы не ожидать «волшебного порядка».

Маленький пример, чтобы увидеть «не по порядку» (разные задержки):

import kotlinx.coroutines.delay
import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.channels.Channel

fun main() = runBlocking {
    val tasks = Channel<Task>()
    val results = Channel<TaskResult>()

    launch {
        for (t in tasks) {
            delay((t.id + 1) * 30L) // разные времена
            results.send(TaskResult(t.id, 0, t.value * 2))
        }
    }

    launch { repeat(3) { tasks.send(Task(it, 10)) }; tasks.close() }

    repeat(3) {
        println(results.receive()) // порядок может отличаться
    }
}

capacity и узкие места

Ёмкость (capacity) — это «сколько задач можно положить в очередь, не дожидаясь, пока их заберут». И это напрямую влияет на поведение producer.

Если capacity = 0 (канал без буфера), send будет ждать, пока кто-то сделает receive. Это похоже на передачу эстафетной палочки: пока второй не взял, первый держит и не может «выложить в коробку».

Если capacity > 0, producer может отправить несколько задач вперёд. Это похоже на ленту на складе: можно положить несколько коробок, и только когда лента заполнится — придётся ждать.

Мини‑эксперимент: один producer быстро шлёт задачи, consumer медленно ест.

import kotlinx.coroutines.delay
import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.channels.Channel

fun main() = runBlocking {
    val tasks = Channel<Int>(capacity = 1)

    val consumer = launch {
        for (x in tasks) {
            delay(100)
            println("consumed $x") // consumed 0, потом consumed 1, ...
        }
    }

    val producer = launch {
        repeat(3) { i ->
            println("sending $i")  // sending 0, sending 1, sending 2
            tasks.send(i)          // на i=2 может подождать, пока consumer освободит место
        }
        tasks.close()
    }

    producer.join()
    consumer.join()
}

Здесь важно не запомнить «волшебные цифры», а понять механику: capacity регулирует, где именно система будет ждать — на producer или на consumer.

6. Небольшой рефакторинг для читаемости

Когда пайплайн становится длиннее, очень хочется всё распилить на функции. Это правильно, но в корутинах легко сделать код менее понятным, если «переусложнить» архитектуру.

Хороший компромисс: вынести создание корутин в небольшие функции‑помощники, которые возвращают Job. Тогда координатору удобно их ждать.

Пример «создать воркера»:

import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Job
import kotlinx.coroutines.delay
import kotlinx.coroutines.launch
import kotlinx.coroutines.channels.ReceiveChannel
import kotlinx.coroutines.channels.SendChannel

fun CoroutineScope.launchWorker(
    workerId: Int,
    tasks: ReceiveChannel<Task>,
    results: SendChannel<TaskResult>
): Job = launch {
    for (t in tasks) {
        delay(50)
        results.send(TaskResult(t.id, workerId, t.value * t.value))
    }
}

И producer:

import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Job
import kotlinx.coroutines.launch
import kotlinx.coroutines.channels.SendChannel

fun CoroutineScope.launchProducer(tasks: SendChannel<Task>, count: Int): Job = launch {
    try {
        repeat(count) { id -> tasks.send(Task(id, id + 1)) }
    } finally {
        // producer — владелец канала задач, он и закрывает
        tasks.close()
    }
}

Такой стиль удобен тем, что main превращается в «сборку схемы», а не в кашу из вложенных launch.

7. Типичные ошибки при Producer–Consumer на Channel

Ошибка №1: забыли закрыть tasks, и воркеры “висят” в for (t in tasks).
Это самый популярный сценарий «почему программа не завершается». Цикл for по каналу заканчивается только тогда, когда канал закрыт и элементы закончились. Если producer завершился без close(), consumers будут честно ждать новые задачи, как сотрудники, которым забыли сказать, что офис переехал. Спасает привычка: закрывать канал в finally у того, кто производит.

Ошибка №2: закрыли results слишком рано.
Здесь логика новичка понятна: «задач больше нет, значит и результатов не будет». Но воркеры могут ещё обрабатывать уже полученные задачи. Если в этот момент закрыть results, следующий results.send(...) у воркера закончится исключением. Правильное закрытие results делается только после join() всех корутин, которые в него отправляют.

Ошибка №3: несколько мест в коде пытаются закрывать один и тот же канал.
Когда close() размазан по коду, вы теряете «владельца протокола» и начинаете ловить странные эффекты: один закрыл, другой ещё работает, третий думает, что канал живой. Лечится договором: у каждого канала должно быть одно логическое место закрытия (producer или координатор), и оно должно быть очевидно при чтении.

Ошибка №4: воркер делает тяжёлую работу, но результат никуда не отправляет, а “печатает сам”.
Для демонстрации это нормально, но для архитектуры — быстро превращается в неудобство. Как только вы захотите вместо печати собирать статистику, считать прогресс или менять формат вывода, окажется, что логика вывода размазана по воркерам. Лучше держать воркера «чистой фабрикой результата»: получил задачу → обработал → отправил результат в results.

Ошибка №5: ожидание “строгого порядка” результатов, как в обычном цикле.
С несколькими воркерами порядок результатов — это порядок готовности, а не порядок задач. Если код где-то неявно рассчитывает «сначала task 0, потом task 1», он начнёт вести себя странно, как только вы добавите второй воркер или измените время обработки. Лучше сразу проектировать вывод как «пришло — обработали», или явно восстанавливать порядок отдельной логикой.

Ошибка №6: игнорирование отмены и отсутствие finally у producer.
Если корутину отменили, или внутри producer случилось исключение, канал задач может остаться открытым, а consumers — зависнуть. try/finally { tasks.close() } делает поведение устойчивым: даже если всё пошло не по плану, система корректно сообщит «новых задач не будет», и consumers завершатся.

1
Задача
Kotlin SELF, 54 уровень, 1 лекция
Недоступна
Лента уведомлений
Лента уведомлений
1
Задача
Kotlin SELF, 54 уровень, 1 лекция
Недоступна
Сумма заказов
Сумма заказов
1
Задача
Kotlin SELF, 54 уровень, 1 лекция
Недоступна
Три работника
Три работника
1
Задача
Kotlin SELF, 54 уровень, 1 лекция
Недоступна
Конвейер результатов
Конвейер результатов
Комментарии
ЧТОБЫ ПОСМОТРЕТЬ ВСЕ КОММЕНТАРИИ ИЛИ ОСТАВИТЬ КОММЕНТАРИЙ,
ПЕРЕЙДИТЕ В ПОЛНУЮ ВЕРСИЮ