Документация

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);
});