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 | |
| Worker (Consumer) | читает задачи из tasks, выполняет обработку, отправляет в results | |
| Collector | читает results и выводит/сохраняет | |
| Coordinator (иногда отдельный) | следит, чтобы каналы закрывались вовремя | |
Важно: воркеров может быть несколько, 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 закрытие каналов — это не «последняя строка программы», а часть контракта.
Логика такая:
- Producer перестал производить задачи → закрываем tasks, чтобы воркеры смогли завершиться.
- Воркеры завершились → они больше не будут отправлять результаты → теперь можно закрыть 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 завершатся.
ПЕРЕЙДИТЕ В ПОЛНУЮ ВЕРСИЮ