Паттерн Идемпотентный Потребитель. Продолжение
Начало
Реализация «Идемпотентного потребителя»
Вот пример идемпотентного потребителя события создания заметки NoteCreated:
internal class NoteCreatedConsumer(
DbContext dbContext,
// … кэш, логгер и т.п.
)
: IConsumer<NoteCreated>
{
public async Task
ConsumeAsync(ConsumeContext<NoteCreated> ctx)
{
// Проверяем, получали ли мы это сообщение
if (await dbCtx
.MessageConsumers.AnyAsync(c =>
c.MessageId == ctx.MessageId &&
c.ConsumerName == nameof(NoteCreatedConsumer)))
return;
using var transaction = await
dbCtx.Database.BeginTransactionAsync();
// … сохраняем заметку в базе
// Записываем, что сообщение обработано
dbCtx.MessageConsumers
.Add(new MessageConsumer
{
MessageId = ctx.MessageId,
ConsumerName = nameof(NoteCreatedConsumer),
ConsumedAtUtc = DateTime.UtcNow
});
await dbContext.SaveChangesAsync();
await transaction.CommitAsync();
// … обновляем кэш, пишем в лог
}
}
Важные детали
1. Ключ Идемпотентности
if (await dbCtx
.MessageConsumers.AnyAsync(c =>
c.MessageId == ctx.MessageId &&
c.ConsumerName == nameof(NoteCreatedConsumer)))
return;
Используем:
- MessageId из контекста передачи сообщений (ctx.MessageId);
- ConsumerName (чтобы несколько получателей могли безопасно обрабатывать одно и то же сообщение).
При поступлении дублирующего сообщения выполняется быстрый выход, и ничего не происходит.
Также важно иметь ограничение уникальности (MessageId, ConsumerName) в таблице MessageConsumers для предотвращения гонок. Таким образом, даже при параллельной обработке одного и того же сообщения только одна сможет вставить запись.
2. Атомарные побочные эффекты + идемпотентность записи
Обработка и сохранение записи получателя сообщения происходят в одной транзакции, поэтому:
- Если обработка завершается неудачей, в таблице MessageConsumers нет записи, поэтому сообщение можно отправить повторно.
- Если обработка завершается успешно, то и заметка, и строка в MessageConsumers фиксируются одновременно.
- Вы никогда не окажетесь в состоянии, когда работа выполнена, но сообщение не помечено как обработанное, и наоборот.
3. Обработка доставки по принципу «как минимум один раз»
Большинство реалистичных конфигураций выполняются по принципу «как минимум один раз»:
- Потребитель обрабатывает сообщение;
- Сбой подтверждения / тайм-аут;
- Брокер повторно доставляет;
- Ваш код выполняется снова.
При использовании этого шаблона второй запуск обращается к таблице MessageConsumers и завершается раньше времени.
Нет дублирования побочных эффектов.
Это работает, за исключением одного нюанса…
Окончание следует…
Источник: https://www.milanjovanovic.tech/blog/the-idempotent-consumer-pattern-in-dotnet-and-why-you-need-it