// Services/PublisherService.csusing System.Threading.Tasks.Dataflow;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
using WorkerService.Dataflow;
using WorkerService.Models;
namespace WorkerService.Services
{
publicclassPublisherService(QuoteChannelchannel, ILogger<PublisherService> logger) : BackgroundService
{
privatereadonlyQuoteChannel _channel = channel;
privatereadonlyILogger<PublisherService> _logger = logger;
privatereadonlyRandom _random = new();
// TPL Dataflow блок для обработки входящих котировок (от Subscriber)privateActionBlock<Quote>? _incomingQuotesBlock;
privateint _publishedCount;
privateint _receivedCount;
protectedoverrideTaskExecuteAsync(CancellationTokenstoppingToken)
{
// Создаем блок для обработки входящих котировок от Subscriber
_incomingQuotesBlock = newActionBlock<Quote>(
quote => ProcessIncomingQuote(quote),
newExecutionDataflowBlockOptions
{
BoundedCapacity = 10,
MaxDegreeOfParallelism = 2,
CancellationToken = stoppingToken,
NameFormat = "Publisher-Incoming"
});
// Подписываемся на обратный канал
_channel.SubscriberToPublisher.LinkTo(
_incomingQuotesBlock,
newDataflowLinkOptions { PropagateCompletion = true });
_logger.LogInformation("Publisher Service запущен. Ожидаем котировки...");
returnPublishQuotesAsync(stoppingToken); // !!!!!!!!!!
}
privateasyncTaskPublishQuotesAsync(CancellationTokenstoppingToken)
{
varsymbols = new[] { "AAPL", "GOOGL", "MSFT", "AMZN", "TSLA" };
while (!stoppingToken.IsCancellationRequested)
{
try
{
varsymbol = symbols[_random.Next(symbols.Length)];
varquote = newQuote
{
Symbol = symbol,
Price = 100 + _random.Next(-20, 20) + (decimal)_random.NextDouble(),
Timestamp = DateTime.UtcNow,
Source = "Publisher"
};
// Публикуем в основной каналawait _channel.PublisherToSubscriber.SendAsync(quote, stoppingToken);
// Отправляем копию в мониторинг
_channel.MonitoringBroadcast.Post(quote);
_publishedCount++;
if (_publishedCount % 10 == 0)
{
_logger.LogInformation("Publisher: опубликовано {Count} котировок", _publishedCount);
}
awaitTask.Delay(TimeSpan.FromMilliseconds(500), stoppingToken);
}
catch (OperationCanceledException)
{
break;
}
catch (Exceptionex)
{
_logger.LogError(ex, "Ошибка при публикации котировки");
}
}
_logger.LogInformation("Publisher завершил публикацию. Всего: {Count}", _publishedCount);
_channel.PublisherToSubscriber.Complete();
}
privatevoidProcessIncomingQuote(Quotequote)
{
_receivedCount++;
_logger.LogInformation("Publisher получил обратную связь: {Quote} (всего получено: {Count})",
quote, _receivedCount);
}
publicoverrideasyncTaskStopAsync(CancellationTokencancellationToken)
{
_logger.LogInformation("Publisher Service останавливается...");
_incomingQuotesBlock?.Complete();
awaitbase.StopAsync(cancellationToken);
}
}
} --------------------------------------------------------------------------------------- Особенно мне понравилось вот это
return PublishQuotesAsync(stoppingToken); - висим на return-е
Шедевр, не правда ли --------------------------------------------------------------------------------------------------------
Отличный код для демонстрации анти-паттерна! Вы абсолютно правы, это действительно «шедевр» в своем роде.
Давайте разберем, почему return PublishQuotesAsync(stoppingToken); — это проблема, и что еще интересного есть в коде.
Главная проблема: return PublishQuotesAsync()
В ExecuteAsync происходит следующее:
Создается ActionBlock и настраивается подписка на канал
Сразу возвращается задача от PublishQuotesAsync
BackgroundService считает, что ExecuteAsync завершился (вернул задачу)
Но PublishQuotesAsync — это бесконечный цикл, который никогда не завершится сам
Почему это плохо?
csharp
// Как должно быть (правильно)protectedoverrideasyncTaskExecuteAsync(CancellationToken stoppingToken){// Инициализация...// await, а не return — тогда метод остается "живым"awaitPublishQuotesAsync(stoppingToken);// Cleanup код здесь выполнится только после завершения PublishQuotesAsync}
Последствия текущей реализации:
BackgroundService думает, что сервис запустился, и переводит его в состояние Running
Если PublishQuotesAsync выбросит исключение, сервис упадет, но StopAsync уже не вызовется корректно
Завершение сервиса через StopAsync сработает, но ExecuteAsync уже вернул управление — это сбивает с толку модель жизненного цикла
Другие интересные моменты
1. Каналы и Dataflow смешаны, но не синхронизированы
csharp
// PublisherToSubscriber — канал (awaitable)await _channel.PublisherToSubscriber.SendAsync(quote, stoppingToken);// MonitoringBroadcast — Dataflow блок (post, не await)
_channel.MonitoringBroadcast.Post(quote);
Мониторинг может терять сообщения, если блок переполнен (BoundedCapacity у MonitoringBroadcast не задан).
2. Нет обработки завершения PublisherToSubscriber
Вы вызываете _channel.PublisherToSubscriber.Complete(), но никто не ждет, пока все сообщения будут обработаны подписчиками.
3. ExecuteAsync не async
Метод не использует async/await, но вызывает асинхронные методы. Это ок, но странно на фоне остального стиля.
4. Гонка при остановке
В StopAsync вы вызываете _incomingQuotesBlock?.Complete(), но PublishQuotesAsync продолжает работать в фоне, пока ее цикл не прервется по stoppingToken. Никто не ждет ее завершения.
Как исправить:
csharp
protectedoverrideasyncTaskExecuteAsync(CancellationToken stoppingToken){
_incomingQuotesBlock =newActionBlock<Quote>(...);
_channel.SubscriberToPublisher.LinkTo(_incomingQuotesBlock,...);
_logger.LogInformation("Publisher Service запущен");try{awaitPublishQuotesAsync(stoppingToken);}finally{
_channel.PublisherToSubscriber.Complete();await _incomingQuotesBlock.Completion;}}
Код демонстрирует хорошее знание TPL Dataflow и Channel, но фундаментальное непонимание жизненного цикла BackgroundService.
Это классический случай, когда код «работает» в тестах (потому что сервис редко останавливают), но в production вылезут баги с зависшими при остановке потоками или необработанными исключениями.
Ваш сарказм по поводу return PublishQuotesAsync() абсолютно заслужен!
Вот совмещённый код — вся логика помещена в метод ExecuteAsync, а вспомогательные фрагменты встроены напрямую:
csharp
protectedoverrideasync Task ExecuteAsync(CancellationToken stoppingToken)
{
// Создаём блок для обработки входящих котировок от Subscriber
_incomingQuotesBlock = new ActionBlock<Quote>(
quote =>
{
_receivedCount++;
_logger.LogInformation(
"Publisher получил обратную связь: {Quote} (всего получено: {Count})",
quote, _receivedCount);
},
new ExecutionDataflowBlockOptions
{
BoundedCapacity = 10,
MaxDegreeOfParallelism = 2,
CancellationToken = stoppingToken,
NameFormat = "Publisher-Incoming"
});
// Подписываемся на обратный канал
_channel.SubscriberToPublisher.LinkTo(
_incomingQuotesBlock,
new DataflowLinkOptions { PropagateCompletion = true });
_logger.LogInformation("Publisher Service запущен. Ожидаем котировки...");
var symbols = new[] { "AAPL", "GOOGL", "MSFT", "AMZN", "TSLA" };
while (!stoppingToken.IsCancellationRequested)
{
try
{
var symbol = symbols[_random.Next(symbols.Length)];
var quote = new Quote
{
Symbol = symbol,
Price = 100 + _random.Next(-20, 20) + (decimal)_random.NextDouble(),
Timestamp = DateTime.UtcNow,
Source = "Publisher"
};
// Публикуем в основной каналawait _channel.PublisherToSubscriber.SendAsync(quote, stoppingToken);
// Отправляем копию в мониторинг
_channel.MonitoringBroadcast.Post(quote);
_publishedCount++;
if (_publishedCount % 10 == 0)
{
_logger.LogInformation("Publisher: опубликовано {Count} котировок", _publishedCount);
}
await Task.Delay(TimeSpan.FromMilliseconds(500), stoppingToken);
}
catch (OperationCanceledException)
{
break;
}
catch (Exception ex)
{
_logger.LogError(ex, "Ошибка при публикации котировки");
}
}
_logger.LogInformation("Publisher завершил публикацию. Всего: {Count}", _publishedCount);
_channel.PublisherToSubscriber.Complete();
}
Что было сделано
Объединение методов. Логика из PublishQuotesAsync полностью перенесена внутрь ExecuteAsync.
Встраивание обработки входящих котировок. Анонимная функция напрямую задана в конструкторе ActionBlock вместо вызова отдельного метода ProcessIncomingQuote. Внутри неё:
увеличивается счётчик _receivedCount;
выполняется логирование получения обратной связи.
Сохранение структуры цикла. Основной цикл генерации и отправки котировок остался без изменений — он теперь выполняется непосредственно в ExecuteAsync.
Сохранение всех обработчиков исключений. Блоки try‑catch для OperationCanceledException и общих исключений сохранены в неизменном виде.
Сохранение логики завершения. После выхода из цикла выполняется:
логирование итогового количества опубликованных котировок;
вызов _channel.PublisherToSubscriber.Complete() для завершения канала отправки котировок.
Сохранение инициализации. Инициализация _incomingQuotesBlock и подписка на канал остались на своих местах в начале метода.
Ключевые особенности итогового кода
Один метод. Вся логика сосредоточена в ExecuteAsync — нет вызовов вспомогательных методов.
Асинхронность. Метод помечен как async Task, что позволяет использовать await внутри цикла.
Обработка отмены. Цикл корректно прерывается при установке stoppingToken.
Логирование. Все ключевые события (запуск, публикация каждых 10‑ти котировок, завершение) логируются.
Обработка ошибок. Исключения корректно перехватываются и логируются, не приводя к аварийному завершению сервиса.
Корректное завершение. Канал отправки котировок завершается вызовом Complete() после выхода из основного цикла.
Как это работает
При запуске сервиса вызывается ExecuteAsync.
Настраивается блок обработки входящих котировок (_incomingQuotesBlock) и устанавливается связь с каналом.
Запускается основной цикл генерации котировок:
генерируется случайная котировка;
отправляется в основной канал и в мониторинг;
обновляется счётчик опубликованных котировок;
каждые 10 котировок выводится логи;
делается пауза 500 мс.
При получении сигнала отмены (stoppingToken) цикл прерывается.
Выполняется финальное логирование и завершается канал отправки котировок.
Хотите, я уточню какой‑то фрагмент кода или предложу дополнительные улучшения?
А вот теперь можно несколько фрагментов инкапсулировать в функции
Например фрагменты бизнес логики инкапсулировать
Алиса
Конечно! Вот код с инкапсуляцией бизнес‑логики в отдельные приватные методы — это улучшит читаемость, поддерживаемость и позволит тестировать отдельные фрагменты:
PublishQuoteAsync — публикует котировку в каналах:
отправляет котировку в основной канал (PublisherToSubscriber);
дублирует котировку в канал мониторинга (MonitoringBroadcast).
UpdateAndLogPublicationStats — обновляет статистику и логирует каждые 10 котировок:
увеличивает счётчик опубликованных котировок;
выводит логи при достижении кратного 10 значения.
Преимущества такого подхода
Читаемость. Код разделён на логические блоки с понятными названиями.
Тестируемость. Каждый метод можно протестировать отдельно (например, проверить генерацию котировки или логику обновления статистики).
Поддерживаемость. Изменения в одной части (например, в алгоритме генерации котировок) не затрагивают другие.
Повторное использование. Методы можно переиспользовать в других сервисах или сценариях.
Соблюдение SRP. Каждый метод выполняет одну задачу, что соответствует принципу единственной ответственности.
Хотите, я раскрою какой‑то конкретный метод подробнее или предложу дополнительные улучшения?
Я
Сейчас стало намного лучше.
Алиса
Спасибо! Рада, что результат вам понравился.
Давайте закрепим понимание — кратко резюмирую, какие именно улучшения появились в коде после инкапсуляции бизнес‑логики в отдельные методы, и почему это действительно делает код лучше.
Что улучшилось
1. Читаемость и структура
Код теперь разделён на чётко очерченные блоки, каждый из которых отвечает за свою задачу. Вместо одного большого метода с множеством строк — несколько маленьких, с говорящими именами. Это позволяет:
быстро понять назначение каждого фрагмента;
легко найти нужный участок кода;
снизить когнитивную нагрузку при чтении.
2. Принцип единственной ответственности (SRP)
Каждый метод выполняет одну конкретную задачу:
SetupIncomingQuotesProcessing — настройка обработки входящих котировок;
GenerateRandomQuote — генерация котировки;
PublishQuoteAsync — публикация котировки в каналах;
UpdateAndLogPublicationStats — обновление статистики и логирование;
ProcessIncomingQuote — обработка обратной связи от подписчиков;
RunQuotePublicationLoop — управление основным циклом публикации.
3. Тестируемость
Теперь можно написать юнит‑тесты для отдельных методов:
Да, стало значительно лучше! Вы правильно исправили главную проблему — теперь ExecuteAsync корректно await-ит цикл публикации. Но есть еще несколько моментов, которые можно доработать. Давайте разберем по порядку.
Вы говорите: "когда SubscriberToPublisher завершится — заверши и _incomingQuotesBlock". Но кто завершает SubscriberToPublisher?
SubscriberService?
Publisher?
Никто?
Риск: Если канал никогда не завершится, блок будет висеть вечно.
2. Нет ожидания завершения Dataflow блоков при остановке
Сейчас в StopAsync:
csharp
publicoverrideasyncTaskStopAsync(CancellationToken cancellationToken){
_logger.LogInformation("Publisher Service останавливается...");
_incomingQuotesBlock?.Complete();// ❌ Никто не ждет Completionawaitbase.StopAsync(cancellationToken);}
Должно быть:
csharp
publicoverrideasyncTaskStopAsync(CancellationToken cancellationToken){
_logger.LogInformation("Publisher Service останавливается...");if(_incomingQuotesBlock !=null){
_incomingQuotesBlock.Complete();await _incomingQuotesBlock.Completion;// Ждем обработки последних сообщений}awaitbase.StopAsync(cancellationToken);}
3. Гонка при доступе к счетчикам
csharp
privatevoidProcessIncomingQuote(Quote quote){
_receivedCount++;// ⚠️ Неатомарная операция в многопоточном контексте}
ActionBlock с MaxDegreeOfParallelism = 2 может вызывать ProcessIncomingQuote из разных потоков. Используйте Interlocked.Increment:
4. PublishQuoteAsync не проверяет успешность отправки мониторингу
csharp
privateasyncTaskPublishQuoteAsync(Quote quote,CancellationToken stoppingToken){await _channel.PublisherToSubscriber.SendAsync(quote, stoppingToken);
_channel.MonitoringBroadcast.Post(quote);// Если блок переполнен - вернет false, вы это игнорируете}
При остановке сначала завершается цикл публикации, потом вызывается Complete() на канале. Но _incomingQuotesBlock продолжает работать и может пытаться отправить логи после завершения. Хорошо бы явно дождаться его завершения в ExecuteAsync:
csharp
protectedoverrideasyncTaskExecuteAsync(CancellationToken stoppingToken){SetupIncomingQuotesProcessing(stoppingToken);
_logger.LogInformation("Publisher Service запущен. Ожидаем котировки...");try{awaitRunQuotePublicationLoop(stoppingToken);}finally{// Даем время дообработать входящие сообщения
_incomingQuotesBlock?.Complete();await(_incomingQuotesBlock?.Completion ?? Task.CompletedTask);}}
Итоговая оценка
С 4/10 (оригинал) вы поднялись до 7/10. Код стал:
Понятным
Поддерживаемым
С правильной асинхронной моделью
Остались "индустриальные" шероховатости (потокобезопасность, graceful shutdown), но это уже детали, которые не убьют прод в первый же день 😄
Главный вопрос, который вы решили правильно:await вместо return в ExecuteAsync. Остальное — полировка. Хорошая работа! 👍
1. Исправлена проблема сPropagateCompletion = true
Было:PropagateCompletion = true— блок автоматически завершался при завершении канала.
Стало:PropagateCompletion = false— теперьPublisherServiceявно управляет завершением блока. Это устраняет риск вечного висения блока, если канал никогда не завершается.
2. Добавлено ожидание завершения Dataflow‑блоков при остановке
ВExecuteAsyncдобавлен блокfinally, где:
вызывается_incomingQuotesBlock?.Complete();
ожидается завершение обработки черезawait (_incomingQuotesBlock?.Completion ?? Task.CompletedTask).
ВStopAsyncтеперь явно ждём завершения_incomingQuotesBlock.Completionпосле вызоваComplete().
Это правильно. Теперь ExecuteAsync явно показывает последовательность: настройка → запуск цикла. В оригинале я возвращал Task от PublishQuotesAsync, что скрывало настройку блока.
3. Устранил лямбду
Было:
_incomingQuotesBlock = new ActionBlock<Quote>(
quote => ProcessIncomingQuote(quote), // Лишняя лямбда
Стало:
_incomingQuotesBlock = new ActionBlock<Quote>(
ProcessIncomingQuote, // Прямая ссылка на метод
...Это чище и эффективнее (нет лишнего аллокационного overhead от лямбды).
Проблема: У тебя MaxDegreeOfParallelism = 2 в _incomingQuotesBlock. Это значит, что два потока одновременно могут вызывать ProcessIncomingQuote, и _receivedCount++ может потерять инкременты.
public override async Task StopAsync(CancellationToken cancellationToken)
{
_logger.LogInformation("Publisher Service останавливается...");
_incomingQuotesBlock?.Complete();
await base.StopAsync(cancellationToken); // ❌ Не ждем завершения!
}
Проблема: Ты вызываешь Complete(), но не ждешь, пока блок обработает все оставшиеся сообщения. Если в буфере есть 5 котировок, они могут быть потеряны.
Решение:
public override async Task StopAsync(CancellationToken cancellationToken)
{
_logger.LogInformation("Publisher Service останавливается...");
_incomingQuotesBlock?.Complete();
// ✅ Ждем, пока все сообщения обработаются
if (_incomingQuotesBlock != null)
{
await _incomingQuotesBlock.Completion;
}
await base.StopAsync(cancellationToken);
}
3. Нет обработки ошибок в ProcessIncomingQuote — СРЕДНЕ
private void ProcessIncomingQuote(Quote quote)
{
_receivedCount++;
_logger.LogInformation(...);
// ❌ Что если здесь выбросится исключение?
}
Проблема: Если в ProcessIncomingQuote (или в логгере) выбросится исключение, блок перейдет в состояние Faulted, и весь конвейер остановится. Ты даже не узнаешь, почему.
Решение:
private void ProcessIncomingQuote(Quote quote)
{
try
{
Interlocked.Increment(ref _receivedCount);
_logger.LogInformation(...);
}
catch (Exception ex)
{
_logger.LogError(ex, "Ошибка при обработке входящей котировки");
// Блок продолжит работу
}
4. MonitoringBroadcast.Post — МЕЛКО, НО ВАЖНО
{
await _channel.PublisherToSubscriber.SendAsync(quote, stoppingToken);
_channel.MonitoringBroadcast.Post(quote); // ❌ Синхронно, может потерять сообщение