Koto.Messaging.Wolverine.Postgres
Надёжное хранилище идемпотентности на PostgreSQL — дедупликация переживает перезапуски.
Зачем
IntegrationEventConsumerBase<TEvent> дедуплицирует события через IProcessedMessageStore. Хранилище in-memory по умолчанию теряет состояние при каждом перезапуске — а с at-least-once доставкой Kafka это означает реальную повторную обработку в продакшене. Этот пакет держит идентификаторы обработанных сообщений в PostgreSQL, так что окно дедупликации переживает передеплой.
Использование
services.AddKotoWolverine();
services.AddPostgresProcessedMessageStore(
builder.Configuration.GetConnectionString("db")!);
Порядок вызова относительно AddKotoWolverine роли не играет.
Опции
services.AddPostgresProcessedMessageStore(connectionString, o =>
{
o.Schema = "koto"; // по умолчанию
o.Table = "processed_messages"; // по умолчанию
o.AutoCreateSchema = true; // по умолчанию; выключите под внешние миграции
o.CleanupInterval = TimeSpan.FromHours(1); // по умолчанию
});
- Окно дедупликации берётся из
KotoWolverineOptions.IdempotencyWindow(по умолчанию 24 ч) — той же опции, что использует in-memory хранилище. - Схема, таблица и индекс создаются при первом обращении (
CREATE ... IF NOT EXISTS). ПоставьтеAutoCreateSchema = false, чтобы управлять ими своими миграциями:
CREATE SCHEMA IF NOT EXISTS koto;
CREATE TABLE IF NOT EXISTS koto.processed_messages (
message_id uuid PRIMARY KEY,
processed_at timestamptz NOT NULL DEFAULT now()
);
CREATE INDEX IF NOT EXISTS ix_processed_messages_processed_at
ON koto.processed_messages (processed_at);
- Фоновый сервис каждые
CleanupIntervalудаляет записи старше окна идемпотентности, так что таблица не растёт бесконечно.
Семантика
Доставка остаётся at-least-once. Хранилище помечает событие как обработанное после того, как завершится ваш ConsumeAsync, вне общей транзакции — падение в промежутке приведёт к повторной доставке. Для строгой бизнес-дедупликации используйте детерминированный operation id с unique-констрейнтом в собственном хранилище консьюмера (см. паттерны ledger в Koto).
Дефолты устойчивого outbox
Та же строка подключения может обслуживать durable outbox Wolverine — конверты в Postgres, EF-транзакции и сбор доменных событий в одном месте:
builder.Host.UseWolverine(opts =>
{
opts.UseKotoKafka(kafka, typeof(SomeHandler).Assembly)
.PublishIntegrationEvents(typeof(OrderPlacedV1).Assembly)
.UseKotoDurableOutbox(connectionString);
});