Паттерн «Производитель-потребитель» c System.Threading.Channels. Продолжение
Начало
Производитель: конечная точка, передающая задачи
Производитель просто выполняет запись. Единственный важный нюанс — как поступить, если канал переполнен; метод
WriteAsync берёт это на себя: он завершается немедленно, если есть свободное место, или ожидает, пока потребитель освободит слот.[ApiController]
[Route("api/[controller]")]
public class ProcessController : ControllerBase
{
private readonly ChannelWriter<WorkItem> _writer;
public ProcessController(Channel<WorkItem> ch)
=> _writer = channel.Writer;
[HttpPost]
public async Task<IActionResult> Enqueue(
WorkItem item, CancellationToken ct)
{
// Ждёт, если канал полный
await _writer.WriteAsync(item, ct);
return Accepted();
}
}
Если при перегрузке системы вы предпочитаете отклонять запросы (что правильнее публичного API, для предотвращения атак), используйте
TryWrite и возвращайте код 429:if (!_writer.TryWrite(item))
return StatusCode(
StatusCodes.Status429TooManyRequests,
"Сервис занят, попробуйте позже.");
return Accepted();
TryWrite никогда не блокирует выполнение. Он возвращает false сразу же, как только канал оказывается заполненным, что позволяет преобразовать эту ситуацию в чёткий ответ 429 вместо ожидания. Выбор стратегии зависит от того, кто инициирует вызов: если это внутренняя пакетная задача, можно позволить ей подождать, а если публичная конечная точка — лучше сразу отклонить запрос.
Потребитель: сервис, считывающий данные из канала
Потребитель реализован в виде BackgroundService, т.е. запускается и останавливается вместе с приложением. Весь цикл обработки умещается в одну строку благодаря методу ReadAllAsync, который выдаёт элементы до тех пор, пока канал не будет закрыт:
public class WorkConsumer : BackgroundService
{
private ChannelReader<WorkItem> _reader;
private ILogger<WorkConsumer> _logger;
public WorkConsumer(
Channel<WorkItem> channel,
ILogger<WorkConsumer> logger)
{
_reader = channel.Reader;
_logger = logger;
}
protected override async Task ExecuteAsync(
CancellationToken stopToken)
{
await foreach (var item in
_reader.ReadAllAsync(stopToken))
{
try
{
await ProcessAsync(item, stopToken);
}
catch (Exception ex)
{
_logger.LogError(ex,
"Ошибка обработки {Id}", item.Id);
}
}
}
private async Task ProcessAsync(
WorkItem item, CancellationToken ct)
{
// … обработка элемента …
}
}
Замечания:
-
try/catch должен находиться внутри цикла и охватывать 1 элемент: если же он будет охватывать весь await foreach, то первое же исключение прервёт цикл, и потребитель завершит работу, в то время как приложение продолжит принимать новые задачи;-
ReadAllAsync корректно завершает выполнение, когда канал закрывается, что обеспечивает штатное завершение работы (об этом далее…).Окончание следует…
Источник: https://thecodeman.net/posts/producer-consumer-with-channels-in-dotnet