TGViewer
.NET Разработчик .NET Разработчик @netdeveloperdiary · 6.75K subscribers
Post #2263 2.58K
День 1873. #ЗаметкиНаПолях
Простая Шина Сообщений в Памяти, Используя Каналы. Окончание

Начало

Потребление событий
Потребление опубликованного IDomainEvent можно реализовать через фоновый сервис (IHostedService). EventProcessorJob использует MessageQueue для чтения (потребления) сообщений. Используем ChannelReader.ReadAllAsync, чтобы получить IAsyncEnumerable, что позволит асинхронно обрабатывать все сообщения в канале.

IPublisher из MediatR поможет связать IDomainEvent с обработчиками. Если мы используем scoped-сервисы, важно получать их из области, для этого внедрим IServiceScopeFactory.
internal class EventProcessorJob(
MessageQueue queue,
IServiceScopeFactory sf,
ILogger<EventProcessorJob> logger)
: BackgroundService
{
protected override async Task
ExecuteAsync(CancellationToken ct)
{
await foreach (IDomainEvent e in
queue.Reader.ReadAllAsync(ct))
{
try
{
using var scope = sf.CreateScope();
var pub = scope.ServiceProvider
.GetRequiredService<IPublisher>();

await pub.Publish(e, ct);
}
catch (Exception ex)
{
logger.LogError(ex,
"Failed! {DomainEventId}",
e.Id);
}
}
}
}

Зарегистрируем сервис:
csharp 
builder.Services
.AddHostedService<EventProcessorJob>();

Использование шины
Сервис IEventBus запишет сообщение в канал и немедленно вернёт управление. Это позволяет публиковать сообщения неблокирующим способом, что повышает производительность. Допустим, при регистрации нового пользователя, нам нужно опубликовать и обработать доменное событие NewUserEvent. Производитель:
internal class UserService(
IUserRepository userRepo,
IEventBus eventBus)
{
public async Task<User> Register(
User user,
CancellationToken ct)
{
// Добавляем пользователя
userRepo.Add(user);

// Публикуем событие
await eventBus.PublishAsync(
new NewUserEvent(user.Id),
ct);

return user;
}
}

Для потребителя нужно определить реализацию INotificationHandler, обрабатывающую доменное событие NewUserEvent, - NewUserEventHandler. Когда фоновое задание EventProcessorJob считает NewUserEvent из канала, оно опубликует сообщение в MediatR и выполнит обработчик.
internal class NewUserEventHandler
: INotificationHandler<NewUserEvent>
{
public async Task Handle(
NewUserEvent event,
CancellationToken ct)
{
// Асинхронно обрабатываем событие, например
// отправим email новому пользователю
}
}

Возможные улучшения
- Устойчивость. Мы можем добавить повторные попытки при возникновении исключений, что повысит надёжность шины сообщений.
- Идемпотентность. Надо ли вы обрабатывать одно и то же сообщение дважды? Если нет, стоит регистрировать обработанные события и проверять перед обработкой, не были ли они обработаны ранее.
- Очередь недоставленных сообщений. Иногда мы не можем правильно обработать сообщение. Можно создать постоянное хранилище для этих сообщений, что позволит устранить неполадки позднее.

Источник: https://www.milanjovanovic.tech/blog/lightweight-in-memory-message-bus-using-dotnet-channels
  • 👍 7
More from @netdeveloperdiary
  1. Oct 8, 2026День 2808. #Карьера 5 Навыков, Которые Помогут Быстрее Стать Сеньором. Начало В ИТ есть се…
  2. Oct 7, 2026День 2807. #ЗаметкиНаПолях Типы Коллекций в .NET, Которые Стоит Попробовать. Окончание Нач…
  3. Oct 6, 2026🦈 Открытое собеседование на Middle C# | 6 октября, 19:00 МСК Приглашаем на открытое собес…
  4. Oct 6, 2026День 2806. #ЗаметкиНаПолях Типы Коллекций в .NET, Которые Стоит Попробовать. Начало Больши…
  5. Oct 5, 2026День 2805. #ЧтоНовенького #NET11 Аргументы в Выражениях Коллекций в C#15 В C#15 реализован…
  6. Oct 4, 2026День 2804. #ВопросыНаСобеседовании Марк Прайс предложил свой набор из 60 вопросов (как тех…
Threads Profile ViewerView any public Threads profile without an account.Open ThreadLook →Writing with AI? Make it sound human.Metric37 rewrites AI drafts so they read naturally. Free AI detector, 1,500 words free.Try Metric37 →