Простая Шина Сообщений в Памяти, Используя Каналы. Окончание
Начало
Потребление событий
Потребление опубликованного 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