Паттерн «Производитель-потребитель» c System.Threading.Channels. Начало
Представим эндпоинт, принимающий файлы, изменяющий их размер и возвращающий ответ. В демо-версии всё работает отлично. Но в проде нагрузка растёт, и эндпойнт получает сотни запросов в секунду. Каждый запрос запускает задачу по изменению размера прямо в потоке обработки; CPU загружается до 100%, а остальные функции приложения начинают завершаться по таймауту, т.к. пул потоков перегружен.
Не обязательно обрабатывать запросы сразу. Достаточно принять данные и обработать их в удобном для системы темпе. Это классическая задача типа «производитель-потребитель»: одна сторона передает задачу, а другая выполняет её со скоростью, которую реально может поддерживать. Для простейшей её реализации не нужны брокеры сообщений, достаточно System.Threading.Channels. Далее рассмотрим реализацию.
Каналы в .NET позволяют передавать данные между производителями и потребителями, работающими параллельно в одном процессе. Производители записывают данные в
ChannelWriter<T>, потребители считывают их из ChannelReader<T>, а канал обеспечивает потокобезопасную асинхронную передачу между ними. Важный момент — канал с ограниченным размером (bounded) автоматически обеспечивает механизм «обратного давления» (backpressure).Почему «наивное» решение только всё усугубляет
Первый порыв в решении проблемы – сделать всё асинхронным: запустить задачу через
Task.Run и сразу вернуть ответ:[HttpPost("process")]
public IActionResult Process(UploadRequest request)
{
// Запустил и забыл. Выглядит асинхронно
_ = Task.Run(() => _imgService.Resize(request));
return Accepted();
}Такой подход обеспечивает быстрый отклик, поэтому кажется удачным решением. Но это не так. Количество запускаемых задач ничем не ограничено, поэтому всплеск нагрузки, который раньше «вешал» CPU, делает это снова — при этом всем отправляется ответ
202, сигнализирующий об успешном принятии запроса. Если происходит перезапуск процесса, эта работа бесследно исчезает. А т.к. за задачами никто не следит, возникающие в них исключения просто пропадают. В итоге вы меняете явное замедление системы на скрытую потерю огромного объёма работы.Правильное решение — использовать очередь с ограничением размера между двумя компонентами. Когда очередь заполняется, отправитель вынужден ждать; это ожидание служит сигналом о том, что задачи поступают быстрее, чем система успевает их обрабатывать, — и именно эту информацию важно выявлять, а не скрывать.
Основы работы с каналами
У
Channel<T> есть два конца. Вы создаёте его и передаёте каждый конец соответствующей стороне:using System.Threading.Channels;
// Вместимость 100. При заполнении писатели ждут
var ch = Channel.CreateBounded<WorkItem>(
new BoundedChannelOptions(100)
{
FullMode = BoundedChannelFullMode.Wait,
SingleReader = false,
SingleWriter = false
});
ChannelWriter<WorkItem> writer = ch.Writer;
ChannelReader<WorkItem> reader = ch.Reader;
Режим переполнения (FullMode) — ключевой элемент:
- Wait —
WriteAsync ждёт появления свободного места (по умолчанию);- DropWrite — молча отбрасывает записываемый элемент;
- DropOldest — заменяет самый старый элемент из очереди входящим;
- DropNewest — заменяет самый новый элемент из очереди входящим.
Режим Wait подходит, когда важен каждый элемент и лучше замедлить работу производителя, чем потерять данные. Режимы Drop* предназначены для телеметрии и потоков данных в реальном времени, где свежий элемент важнее полной истории: например, при передаче метрик отбросить самое старое показание вполне допустимо.
Параметры SingleReader и SingleWriter служат для оптимизации. Устанавливайте их в
true, только когда у вас действительно ровно один читатель или один писатель — так канал использует более быстрый внутренний путь обработки. Если сомневаетесь, оставляйте false: ошибочный true может привести к трудноуловимому состоянию гонки.Запросы поступают из множества потоков и записываются в один канал. Опустошением канала занимается небольшой фиксированный пул потребителей. Когда канал переполняется, операция записи приостанавливается, и сигнал обратного давления передаётся непосредственно вызывающему коду — именно это и требуется, так как замедление становится явным и предсказуемым.
Продолжение следует…
Источник: https://thecodeman.net/posts/producer-consumer-with-channels-in-dotnet