Простая Шина Сообщений в Памяти, Используя Каналы. Начало
Обмен сообщениями играет важную роль в современной архитектуре ПО, обеспечивая координацию между слабосвязанными компонентами. Шина сообщений в памяти особенно полезна, когда критическими требованиями являются высокая производительность и низкая задержка.
Недостатки
- Потеря сообщений при сбое процесса программы.
- Работает только внутри одного процесса, поэтому бесполезна в распределённых системах.
Практический вариант использования — создание модульного монолита. Можно реализовать связь между модулями через доменные события. Когда нужно будет выделить какой-то модуль в отдельный сервис, можно заменить шину в памяти на распределённую.
Абстракции для шины
Нам нужны две абстракции для публикации сообщений и для обработчика сообщений.
Интерфейс IEventBus предоставляет метод PublishAsync для публикации сообщений. Также определено ограничение, которое позволяет передавать только экземпляр IDomainEvent.
public interface IEventBus
{
Task PublishAsync<T>(
T domainEvent,
CancellationToken ct = default)
where T : class, IDomainEvent;
}
Используем MediatR для модели издатель-подписчик. Интерфейс IDomainEvent будет наследоваться от INotification. Это позволит легко определять обработчики IDomainEvent с помощью INotificationHandler<T>. Кроме того, в IDomainEvent добавим идентификатор, чтобы отслеживать выполнение. Абстрактный класс DomainEvent будет базовым для конкретных реализаций.
using MediatR;
public interface IDomainEvent : INotification
{
Guid Id { get; init; }
}
public abstract record DomainEvent(Guid Id)
: IDomainEvent;
Простая очередь в памяти с использованием каналов
Пространство имён System.Threading.Channels предоставляет структуры данных для асинхронной передачи сообщений между производителями и потребителями. Производители асинхронно создают данные, а потребители асинхронно потребляют их. В отличие от традиционных очередей сообщений, каналы полностью работают в памяти. Недостатком этого подхода является возможность потери сообщения в случае сбоя приложения.
MessageQueue создаёт неограниченный канал, т.е. у канала может быть любое количество читателей и писателей. Он также предоставляет ChannelReader и ChannelWriter, которые позволяют клиентам публиковать и потреблять сообщения.
internal class MessageQueue
{
private readonly Channel<IDomainEvent> _channel =
Channel.CreateUnbounded<IDomainEvent>();
public ChannelReader<IDomainEvent>
Reader => _channel.Reader;
public ChannelWriter<IDomainEvent>
Writer => _channel.Writer;
}
Зарегистрируем сервис как синглтон:
builder.Services.AddSingleton<MessageQueue>();
Реализация шины событий
Класс EventBus использует MessageQueue для доступа к ChannelWriter и записи события в канал.
internal class EventBus(MessageQueue queue)
: IEventBus
{
public async Task PublishAsync<T>(
T domainEvent,
CancellationToken ct = default)
where T : class, IDomainEvent
{
await queue.Writer.WriteAsync(
domainEvent, ct);
}
}
Шину также регистрируем как синглтон, т.к. она не сохраняет состояния:
builder.Services.AddSingleton<IEventBus, EventBus>();
Окончание следует…
Источник: https://www.milanjovanovic.tech/blog/lightweight-in-memory-message-bus-using-dotnet-channels