Привет, Хабр!
Задача «одна часть кода производит данные, другая их обрабатывает» встречается примерно в каждом сервисе.
Приняли запрос — положили в очередь на обработку.
Прочитали файл — отдали парсеру.
Собрали события — отправили пачкой в аналитику.
Исторически это писали на ConcurrentQueue с ручным сигналом, на BlockingCollection с блокирующими потоками или на самодельной обёртке с SemaphoreSlim. Всё это работает, но каждая такая конструкция норовит обзавестись собственными багами: то потребитель крутится вхолостую, то поток блокируется там, где не должен, то очередь молча растёт, пока не кончится память.
System.Threading.Channels закрывает задачу целиком: асинхронное чтение и запись, ограничение размера с обратным давлением, корректное завершение. Библиотека небольшая, API компактный, и именно поэтому её часто используют вполсилы — берут CreateUnbounded, пишут await foreach и считают, что готово.
А потом сервис молча отъедает память, воркер не выключается при остановке приложения, а первое же исключение внутри обработчика убивает конвейер без единой записи в логах. Разберём, как этого не допустить.
Из чего состоит канал
Канал — это очередь с двумя половинами: Writer для producer и Reader для consumer. Разделение не косметическое: вы отдаёте в разные части кода разные интерфейсы, и потребитель физически не может записать в канал, а производитель — прочитать.
var channel = Channel.CreateUnbounded<WorkItem>(); ChannelWriter<WorkItem> writer = channel.Writer; ChannelReader<WorkItem> reader = channel.Reader;
Простейший конвейер выглядит так:
// производитель _ = Task.Run(async () => { for (int i = 0; i < 100; i++) await writer.WriteAsync(new WorkItem(i)); writer.Complete(); }); // потребитель await foreach (var item in reader.ReadAllAsync()) { await ProcessAsync(item); }
ReadAllAsync возвращает асинхронную последовательность, которая сама завершится, когда канал закроют и разберут остаток.
Эти пятнадцать строк — вся базовая механика. Дальше про то, из‑за чего конвейеры могут ломаться.
Unbounded — это не «без ограничений», это «ограничение в виде OOM»
Первая и самая частая ошибка — взять неограниченный канал по умолчанию, потому что так проще.
var channel = Channel.CreateUnbounded<Order>();
Пока производитель медленнее потребителя, всё прекрасно. Проблема в том, что «медленнее» — это не свойство системы, а стечение обстоятельств. База притормозила, внешний API стал отвечать за две секунды вместо ста миллисекунд, приехал всплеск трафика — и вот производитель уже быстрее.
Очередь начинает расти. Она будет расти, пока не кончится память, и никакого сигнала об этом вы не получите: метрики приложения в порядке, ошибок нет, запросы принимаются. Потом процесс падает по OOM, и в логах ничего полезного.
Ограниченный канал превращает эту ситуацию из аварии в штатное замедление:
var channel = Channel.CreateBounded<Order>(new BoundedChannelOptions(1000) { FullMode = BoundedChannelFullMode.Wait, });
Теперь при заполнении WriteAsync не завершается, пока не освободится место. Производитель притормаживает ровно настолько, насколько отстал потребитель.
Обратное давление должно доходить до источника. Если ваш производитель — это HTTP‑обработчик, то await writer.WriteAsync(...) заставит подождать сам запрос, и клиент увидит замедление.
Что делать, когда канал полон
Не всегда правильный ответ — ждать. BoundedChannelFullMode даёт четыре варианта поведения, и выбор между ними — это решение о том, что для вас важнее.
Wait— производитель ждёт места. Подходит, когда терять данные нельзя: заказы, платежи, задачи на обработку.DropOldest— самый старый элемент выбрасывается, новый занимает его место. Логичный выбор для телеметрии и текущего состояния: свежие показания датчика ценнее позавчерашних.DropNewest— выбрасывается последний добавленный, новый тоже не проходит.DropWrite— новый элемент просто не попадает в канал, очередь остаётся как была.
Отдельно стоит знать про перегрузку с колбэком на выброшенный элемент:
var channel = Channel.CreateBounded<Metric>( new BoundedChannelOptions(500) { FullMode = BoundedChannelFullMode.DropOldest }, itemDropped: m => _droppedCounter.Add(1));
Без этого счётчика режимы с отбрасыванием опасны.
Есть ещё вариант «попробовать записать и не ждать»:
if (!writer.TryWrite(item)) { _logger.LogWarning("Очередь переполнена, элемент отклонён"); return Results.StatusCode(503); }
TryWrite возвращает false, если места нет. Для веб‑сервиса это часто самый разумный вариант: лучше честно ответить «перегружены, попробуйте позже», чем держать соединение открытым неизвестно сколько.
Завершение: место, где конвейеры зависают
Второй по частоте баг — забытый Complete().
foreach (var item in source) await writer.WriteAsync(item); // Complete() не вызвали
Потребитель в await foreach будет ждать вечно. Приложение не упадёт, ошибок не будет, просто фоновая задача никогда не закончится и процесс не завершится штатно. Отлаживать такое неприятно: симптом — «сервис не останавливается», а причина в другом файле.
Правильный шаблон — try/finally:
try { foreach (var item in source) await writer.WriteAsync(item, ct); } finally { writer.Complete(); }
Теперь канал закроется в любом случае, даже если производитель упал с исключением.
Зеркальная ошибка — писать в закрытый канал. После Complete() любой WriteAsync бросит ChannelClosedException. Если у вас несколько производителей и один из них может закончить раньше остальных, Complete() вызывает не он, а координирующий код:
var producers = Enumerable.Range(0, 4) .Select(i => ProduceAsync(writer, i, ct)) .ToArray(); await Task.WhenAll(producers); writer.Complete(); // закрываем, когда закончили все
Проверить, закрыт ли канал, можно через TryComplete() — он возвращает false вместо исключения, если канал уже закрыт.
Ошибки: как не потерять исключение
Допустим, производитель упал:
try { await foreach (var row in ReadFromDatabaseAsync(ct)) await writer.WriteAsync(row, ct); } finally { writer.Complete(); // канал закрыт, потребитель спокойно завершился }
Потребитель увидит нормальное завершение и решит, что данные кончились. На самом деле они не кончились — произошла ошибка, но никто об этом не узнал. В базу уедет половина данных, и всё будет выглядеть штатно.
Complete() умеет принимать исключение:
try { await foreach (var row in ReadFromDatabaseAsync(ct)) await writer.WriteAsync(row, ct); writer.Complete(); } catch (Exception ex) { writer.Complete(ex); // потребитель получит это исключение }
Теперь ReadAllAsync у потребителя бросит переданное исключение после того, как отдаст все успевшие записаться элементы. Ошибка доходит до того, кто может её обработать, а не растворяется в фоновой задаче.
С обратной стороны действует то же правило. Если исключение вылетит из тела await foreach, цикл прервётся, а канал останется открытым — производитель продолжит писать в никуда, пока не упрётся в границу и не зависнет навсегда.
Поэтому обработчик защищают внутри:
await foreach (var item in reader.ReadAllAsync(ct)) { try { await ProcessAsync(item, ct); } catch (Exception ex) { _logger.LogError(ex, "Не смогли обработать {Id}", item.Id); // решаем: пропустить, отправить в очередь ошибок, повторить } }
Одно битое сообщение не должно валить весь конвейер, но и молча глотать все ошибки подряд тоже плохо. Обычно правильный ответ такой: ошибку обработки конкретного элемента логируем и идём дальше, а вот исключения инфраструктурного уровня (потеряли соединение с базой) пробрасываем наружу, чтобы конвейер перезапустился целиком.
Несколько потребителей
Один потребитель — узкое место, если обработка тяжёлая. Каналы поддерживают любое количество читателей, и параллелизм получается почти бесплатно:
var consumers = Enumerable.Range(0, Environment.ProcessorCount) .Select(_ => Task.Run(async () => { await foreach (var item in reader.ReadAllAsync(ct)) await ProcessAsync(item, ct); })) .ToArray(); await Task.WhenAll(consumers);
Каждый элемент достанется ровно одному потребителю — канал сам разбирается с конкуренцией. Никаких блокировок писать не нужно.
Единственное, о чём стоит помнить: порядок обработки при этом теряется. Если элементы связаны между собой (события одного пользователя, операции по одному счёту), параллельная обработка их перемешает. Решается либо шардированием — по каналу на ключ, либо сохранением одного потребителя там, где порядок важен.
И раз уж речь о нескольких читателях: если их точно один, скажите об этом каналу.
var channel = Channel.CreateBounded<Order>(new BoundedChannelOptions(1000) { SingleReader = true, SingleWriter = false, });
Флаги SingleReader и SingleWriter позволяют реализации использовать более дешёвые внутренние структуры без части синхронизации. Выигрыш заметен на высокой частоте, но соврать тут нельзя: поставите SingleReader = true и запустите два потребителя, получите повреждённое состояние без всяких предупреждений.
Есть ещё AllowSynchronousContinuations. Он разрешает выполнять продолжение прямо в потоке того, кто записал элемент, экономя переключение контекста. Уменьшает задержку, но при неудачном стечении обстоятельств производитель начинает выполнять работу потребителя. По умолчанию выключен, и в большинстве случаев так и надо оставить.
Чем каналы лучше того, что было раньше
Раз уж дошли до флагов производительности, стоит объяснить, почему вообще стоит переезжать с привычных конструкций.
BlockingCollection<T>— самый близкий аналог. Проблема в том, что он блокирующий: потребитель, ожидающий элемент, занимает поток целиком. Десять потребителей — десять занятых потоков пула, которые ничего не делают. В асинхронном приложении это прямой путь к голоданию пула, когда свободных потоков не остаётся и встаёт вообще всё.ConcurrentQueue<T>не блокирует, но и не умеет ждать. Приходится либо крутить опрос в цикле с задержкой, либо прикручиватьSemaphoreSlimдля сигнализации — и вот вы уже пишете свою реализацию канала, только с багами.
Канал не занимает поток на ожидании: await возвращает управление, и поток уходит делать другую работу. Плюс внутри применяются оптимизации, которых у самодельного решения не будет: возврат ValueTask вместо Task там, где операция завершается синхронно, отсутствие лишних аллокаций на каждый элемент, разные внутренние реализации под разные комбинации флагов.
Практически это значит, что канал с сотней потребителей не съест сотню потоков — он вообще не потратит ни одного, пока данных нет.
Отмена, сделанная правильно
Токен отмены передаётся и в запись, и в чтение — и делает он в этих местах разные вещи.
await writer.WriteAsync(item, ct); // отменит ожидание места в полном канале await foreach (var item in reader.ReadAllAsync(ct)) // отменит ожидание новых элементов
Отмена чтения не закрывает канал. Потребитель просто перестанет ждать и выйдет из цикла с OperationCanceledException, а канал останется жив со всем содержимым. Если вам нужно именно закрыть его при остановке, делайте это явно.
Для корректного завершения при остановке приложения обычно нужно перестать принимать новое, доработать то, что уже в очереди.
public async Task StopAsync(CancellationToken ct) { _writer.Complete(); // больше ничего не принимаем await _consumerTask.WaitAsync(ct); // ждём, пока разберут остаток }
WaitAsync тут не даст висеть вечно, если потребитель по какой‑то причине застрял.
Каналы в ASP.NET Core
Самый частый сценарий — фоновая обработка запросов. Собирается это из канала и BackgroundService.
public sealed class WorkQueue { private readonly Channel<WorkItem> _channel = Channel.CreateBounded<WorkItem>( new BoundedChannelOptions(10_000) { FullMode = BoundedChannelFullMode.Wait }); public ChannelReader<WorkItem> Reader => _channel.Reader; public bool TryEnqueue(WorkItem item) => _channel.Writer.TryWrite(item); public void Complete() => _channel.Writer.Complete(); }
Обёртка нужна не ради красоты: наружу торчит только то, что можно, а TryEnqueue не даёт контроллеру случайно заблокироваться на полной очереди.
public sealed class WorkProcessor(WorkQueue queue, IServiceScopeFactory scopes) : BackgroundService { protected override async Task ExecuteAsync(CancellationToken stoppingToken) { await foreach (var item in queue.Reader.ReadAllAsync(stoppingToken)) { using var scope = scopes.CreateScope(); var handler = scope.ServiceProvider.GetRequiredService<IWorkHandler>(); try { await handler.HandleAsync(item, stoppingToken); } catch (Exception ex) when (ex is not OperationCanceledException) { // логируем и продолжаем } } } }
Обратите внимание на область видимости. BackgroundService регистрируется как синглтон, а большинство сервисов в ASP.NET живут в области запроса, тот же DbContext.
Поэтому scope создаётся на каждый элемент вручную. Если этого не сделать и внедрить DbContext прямо в конструктор фонового сервиса, вы получите один контекст на всё время жизни приложения: он накопит отслеживаемые сущности, съест память и рано или поздно упадёт на конкурентном доступе.
И фильтр when (ex is not OperationCanceledException) тут не для красоты. При остановке приложения токен сработает, и отмена не должна попасть в лог как ошибка.
Если хотите проверить, насколько уверенно вы работаете с асинхронностью, управлением памятью и устройством.NET, пройдите вступительное тестирование по C#. Оно покажет текущий уровень и темы, которые стоит разобрать глубже.
Тестировать это тоже надо
Конвейеры неприятно тестировать, потому что в них замешана асинхронность и время.
Не тестируйте канал, тестируйте обработчик. Логика обработки одного элемента — обычный метод, он проверяется обычным юнит‑тестом без всяких каналов.
Для проверки самого конвейера удобно писать в канал вручную и собирать результат:
[Fact] public async Task Processes_all_items() { var channel = Channel.CreateUnbounded<int>(); var processed = new List<int>(); var consumer = Task.Run(async () => { await foreach (var i in channel.Reader.ReadAllAsync()) processed.Add(i); }); foreach (var i in Enumerable.Range(0, 100)) await channel.Writer.WriteAsync(i); channel.Writer.Complete(); await consumer.WaitAsync(TimeSpan.FromSeconds(5)); Assert.Equal(100, processed.Count); }
WaitAsync с таймаутом тут обязателен.
К тому же сценарий с переполнением проверяется каналом ёмкостью в один‑два элемента. На тысяче элементов вы обратное давление не воспроизведёте, а на двух легко.
Конвейер из нескольких стадий
Каналы хорошо соединяются друг с другом, и из них получается многоступенчатая обработка, где каждая стадия работает со своей скоростью.
static ChannelReader<TOut> Stage<TIn, TOut>( ChannelReader<TIn> input, Func<TIn, CancellationToken, ValueTask<TOut>> transform, int capacity, int parallelism, CancellationToken ct) { var output = Channel.CreateBounded<TOut>(capacity); var workers = Enumerable.Range(0, parallelism).Select(_ => Task.Run(async () => { await foreach (var item in input.ReadAllAsync(ct)) await output.Writer.WriteAsync(await transform(item, ct), ct); }, ct)); _ = Task.WhenAll(workers).ContinueWith(t => output.Writer.Complete(t.Exception), TaskScheduler.Default); return output.Reader; }
Собирается это в цепочку:
var parsed = Stage(rawReader, ParseAsync, capacity: 500, parallelism: 4, ct); var enriched = Stage(parsed, EnrichAsync, capacity: 200, parallelism: 8, ct); var ready = Stage(enriched, ValidateAsync, capacity: 100, parallelism: 2, ct);
Красота такой схемы в том, что каждая стадия имеет свою пропускную способность и свой буфер. Обогащение ходит по сети и требует восьми параллельных обработчиков, валидация упирается в процессор и обходится двумя. Обратное давление распространяется по всей цепочке само: забился последний буфер — притормозит первый.
Обратите внимание на Complete(t.Exception) в продолжении, так исключение из любой стадии доедет до следующей, а не потеряется в фоновой задаче.
Пакетирование и приоритеты
Писать в базу по одной строке — расточительство, а копить пачку вручную неудобно.
static async IAsyncEnumerable<List<T>> Batch<T>( ChannelReader<T> reader, int maxSize, TimeSpan maxDelay, [EnumeratorCancellation] CancellationToken ct = default) { var batch = new List<T>(maxSize); using var timer = new PeriodicTimer(maxDelay); while (await reader.WaitToReadAsync(ct)) { while (batch.Count < maxSize && reader.TryRead(out var item)) batch.Add(item); if (batch.Count > 0) { yield return batch; batch = new List<T>(maxSize); } } }
Основные методы здесь — WaitToReadAsync и TryRead. Первый ждёт, пока в канале появится хоть что‑то, второй забирает элементы без ожидания, пока они есть. Вместе они дают то поведение, которое нужно для пакетов: дождались первого элемента, быстро выгребли всё, что накопилось, отдали пачку.
Ограничение по времени тут не менее важно, чем по размеру. Без него последние девять элементов из десяти будут лежать в буфере, пока не приедет десятый.
А еще есть всем известная возможность: с.NET 9 появились приоритетные каналы.
var channel = Channel.CreateUnboundedPrioritized<Job>( new UnboundedPrioritizedChannelOptions<Job> { Comparer = Comparer<Job>.Create((a, b) => a.Priority.CompareTo(b.Priority)), });
Читаться будет элемент с наименьшим значением по компаратору.
Где каналы не нужны
Когда данные нельзя терять при перезапуске. Канал живёт в памяти процесса. Упало приложение — очередь исчезла вместе с ним. Если задачи нужно гарантированно выполнить, нужен настоящий брокер с сохранением на диск: RabbitMQ, Kafka, очередь в базе. Канал внутри процесса при этом остаётся полезным как буфер между брокером и обработчиком.
Когда нужно распределить работу между экземплярами. Канал не выходит за границы процесса. Три пода в кластере — три независимые очереди, каждая со своей нагрузкой.
Когда данных мало. Десять элементов раз в минуту не требуют конвейера, обратного давления и параллельных потребителей. Обычный вызов метода будет понятнее и короче.
Когда нужны сложные преобразования потока. Слияние, ветвление, оконные операции, повторы с задержкой — всё это можно собрать на каналах руками, но библиотеки реактивных расширений делают это выразительнее. Каналы сильны в простой и предсказуемой передаче данных, а не в декларативной обработке потоков.

После работы с каналами быстро выясняется, что одних знаний синтаксиса C# недостаточно: важны архитектура, асинхронность и понимание поведения приложения под нагрузкой.
На бесплатных занятиях можно разобрать эти темы с практиками, задать вопросы и посмотреть, как устроено обучение.
5 августа в 20:00. «Как AI меняет работу C#‑разработчика». Записаться
18 августа в 20:00. «Архитектурные ошибки, которые совершают даже опытные C#‑разработчики». Записаться
Больше бесплатных уроков и других полезных подборок смотрите в дайджесте.
