Архитектура

Паттерны управления состоянием в распределённых системах: Saga, Outbox и Event Sourcing с PostgreSQL

Ruslan Ismailov Опубликовано 14 мин чтения
П

Введение: проблема распределённых транзакций

Когда монолит разбивается на микросервисы, первой жертвой становится транзакционность. В классическом ACID-мире одна база данных гарантирует атомарность: либо всё, либо ничего. В распределённой системе та же операция — например, оформление заказа — затрагивает сервис заказов, сервис оплаты, сервис склада и сервис уведомлений. Каждый из них имеет собственную базу данных, и единой транзакции между ними не существует.

Протокол двухфазной фиксации (2PC) теоретически решает эту проблему, но на практике создаёт узкое место: координатор может упасть в момент между фазами, оставив систему в неопределённом состоянии. Кроме того, 2PC блокирует ресурсы участников на время всей транзакции, что убивает производительность при высокой нагрузке.

Архитектурное сообщество ответило на этот вызов тремя паттернами, которые стали стандартом де-факто в 2024–2026 годах: Saga, Transactional Outbox и Event Sourcing. В этой статье мы детально разберём каждый из них, а затем покажем практическую реализацию на Go с PostgreSQL и Redis.

Паттерн Saga: управление долгоживущими транзакциями

Saga — это последовательность локальных транзакций, каждая из которых публикует событие или сообщение, запускающее следующий шаг. Если один шаг завершается неудачей, Saga выполняет компенсирующие транзакции для всех предшествующих шагов, откатывая систему в согласованное состояние.

Хореография vs Оркестрация

Существуют два варианта реализации Saga, кардинально отличающихся по архитектуре:

  • Хореография (Choreography): каждый сервис реагирует на события и публикует собственные. Нет центрального координатора. Сервис заказов публикует OrderCreated, сервис оплаты подписывается на это событие, резервирует средства и публикует PaymentReserved, сервис склада подписывается на PaymentReserved и так далее.
  • Оркестрация (Orchestration): центральный оркестратор (Saga Orchestrator) явно командует каждым сервисом: «заблокируй оплату», «зарезервируй товар», «отправь уведомление». Оркестратор отслеживает состояние всего процесса.

Хореография проще в начале — нет центральной точки отказа — но по мере роста числа сервисов граф зависимостей становится трудноотлаживаемым. Оркестрация даёт централизованный контроль и наблюдаемость, но оркестратор превращается в точку концентрации бизнес-логики и потенциальное узкое место.

Компенсирующие транзакции

Ключевое требование к Saga: каждая локальная транзакция должна иметь компенсирующую пару. Если шаг «списать деньги» прошёл успешно, а следующий шаг «зарезервировать товар» завершился ошибкой, система обязана выполнить «вернуть деньги». Компенсации должны быть идемпотентными: повторный вызов не должен приводить к двойному возврату.

Плюсы Saga: нет распределённых блокировок, высокая доступность, каждый сервис независим. Минусы: временная несогласованность (eventual consistency), сложность проектирования компенсаций, трудности с отладкой в хореографии.

Transactional Outbox Pattern: гарантированная доставка событий

Saga и Event Sourcing предполагают надёжную публикацию событий. Но как гарантировать, что событие будет опубликовано, даже если сервис упал сразу после коммита транзакции? Именно здесь на сцену выходит Transactional Outbox.

Идея проста и элегантна: вместо того чтобы публиковать событие напрямую в брокер сообщений, сервис записывает событие в специальную таблицу outbox в той же транзакции, что и бизнес-данные. Отдельный фоновый процесс (Message Relay или Polling Worker) читает непрочитанные записи из outbox и публикует их в брокер. Атомарность гарантируется базой данных: либо и бизнес-запись, и outbox-запись созданы, либо ни одна.

SQL-схема для Outbox в PostgreSQL

CREATE TABLE outbox_events (
  id          UUID PRIMARY KEY DEFAULT gen_random_uuid(),
  aggregate_type VARCHAR(100) NOT NULL,
  aggregate_id   VARCHAR(100) NOT NULL,
  event_type     VARCHAR(100) NOT NULL,
  payload        JSONB        NOT NULL,
  created_at     TIMESTAMPTZ  NOT NULL DEFAULT NOW(),
  published_at   TIMESTAMPTZ,
  retry_count    INT          NOT NULL DEFAULT 0,
  status         VARCHAR(20)  NOT NULL DEFAULT 'PENDING'
    CHECK (status IN ('PENDING', 'PUBLISHED', 'FAILED'))
);

CREATE INDEX idx_outbox_status_created
  ON outbox_events (status, created_at)
  WHERE status = 'PENDING';

Индекс по частичному условию (WHERE status = 'PENDING') критически важен для производительности: polling worker будет обращаться к этому индексу при каждой итерации, и без него полный скан таблицы станет проблемой уже при нескольких миллионах записей.

Event Sourcing: состояние как последовательность событий

Event Sourcing переворачивает традиционную модель хранения данных. Вместо того чтобы хранить текущее состояние сущности (обновляя строку в таблице), мы храним последовательность событий, которые привели к этому состоянию. Текущее состояние восстанавливается путём воспроизведения (replay) всех событий.

Например, для банковского счёта не хранится поле balance = 1500. Вместо этого хранятся события: AccountOpened(0), MoneyDeposited(2000), MoneyWithdrawn(500). Баланс 1500 вычисляется при чтении.

Хранение событий в PostgreSQL

CREATE TABLE event_store (
  id             BIGSERIAL    PRIMARY KEY,
  stream_id      VARCHAR(200) NOT NULL,
  stream_version BIGINT       NOT NULL,
  event_type     VARCHAR(100) NOT NULL,
  event_data     JSONB        NOT NULL,
  metadata       JSONB        NOT NULL DEFAULT '{}',
  created_at     TIMESTAMPTZ  NOT NULL DEFAULT NOW(),
  UNIQUE (stream_id, stream_version)
);

CREATE INDEX idx_event_store_stream
  ON event_store (stream_id, stream_version);

Ограничение UNIQUE (stream_id, stream_version) обеспечивает оптимистичную блокировку: если два процесса попытаются записать событие с одинаковой версией для одного стрима, PostgreSQL отклонит одну из операций, предотвращая конфликт.

Снапшоты (Snapshots): при больших стримах полный replay становится дорогостоящим. Решение — периодически сохранять снапшот текущего состояния и воспроизводить только события, произошедшие после снапшота.

Практическая реализация Outbox на Go с PostgreSQL

Рассмотрим полную реализацию Polling Worker для Transactional Outbox на Go. Используем pgx для работы с PostgreSQL и стандартные паттерны конкурентности.

Структуры данных

package outbox

import (
    "context"
    "time"
    "github.com/google/uuid"
)

type Event struct {
    ID            uuid.UUID
    AggregateType string
    AggregateID   string
    EventType     string
    Payload       []byte
    CreatedAt     time.Time
    RetryCount    int
}

type Publisher interface {
    Publish(ctx context.Context, event Event) error
}

Polling Worker

package outbox

import (
    "context"
    "fmt"
    "log/slog"
    "time"

    "github.com/jackc/pgx/v5/pgxpool"
)

type Worker struct {
    db        *pgxpool.Pool
    publisher Publisher
    batchSize int
    interval  time.Duration
}

func NewWorker(db *pgxpool.Pool, pub Publisher) *Worker {
    return &Worker{
        db:        db,
        publisher: pub,
        batchSize: 100,
        interval:  500 * time.Millisecond,
    }
}

func (w *Worker) Run(ctx context.Context) error {
    ticker := time.NewTicker(w.interval)
    defer ticker.Stop()

    for {
        select {
        case <-ctx.Done():
            return ctx.Err()
        case <-ticker.C:
            if err := w.processBatch(ctx); err != nil {
                slog.Error("outbox: batch processing failed", "error", err)
            }
        }
    }
}

func (w *Worker) processBatch(ctx context.Context) error {
    tx, err := w.db.Begin(ctx)
    if err != nil {
        return fmt.Errorf("begin tx: %w", err)
    }
    defer tx.Rollback(ctx)

    // SELECT FOR UPDATE SKIP LOCKED — ключ к горизонтальному масштабированию
    rows, err := tx.Query(ctx, `
        SELECT id, aggregate_type, aggregate_id, event_type, payload, created_at, retry_count
        FROM outbox_events
        WHERE status = 'PENDING'
        ORDER BY created_at
        LIMIT $1
        FOR UPDATE SKIP LOCKED
    `, w.batchSize)
    if err != nil {
        return fmt.Errorf("query events: %w", err)
    }

    var events []Event
    for rows.Next() {
        var e Event
        if err := rows.Scan(
            &e.ID, &e.AggregateType, &e.AggregateID,
            &e.EventType, &e.Payload, &e.CreatedAt, &e.RetryCount,
        ); err != nil {
            return fmt.Errorf("scan event: %w", err)
        }
        events = append(events, e)
    }
    rows.Close()

    for _, event := range events {
        if err := w.publisher.Publish(ctx, event); err != nil {
            slog.Warn("outbox: publish failed", "event_id", event.ID, "error", err)
            _, _ = tx.Exec(ctx,
                `UPDATE outbox_events SET retry_count = retry_count + 1,
                 status = CASE WHEN retry_count >= 4 THEN 'FAILED' ELSE 'PENDING' END
                 WHERE id = $1`, event.ID)
            continue
        }
        _, _ = tx.Exec(ctx,
            `UPDATE outbox_events SET status = 'PUBLISHED', published_at = NOW() WHERE id = $1`,
            event.ID)
    }

    return tx.Commit(ctx)
}

Директива FOR UPDATE SKIP LOCKED в PostgreSQL позволяет запускать несколько экземпляров polling worker параллельно: каждый захватит свой набор строк, не конкурируя с остальными. Это обеспечивает горизонтальное масштабирование обработки Outbox.

Запись в Outbox внутри бизнес-транзакции

func (s *OrderService) CreateOrder(ctx context.Context, req CreateOrderRequest) error {
    tx, err := s.db.Begin(ctx)
    if err != nil {
        return err
    }
    defer tx.Rollback(ctx)

    // 1. Создаём заказ
    orderID := uuid.New()
    _, err = tx.Exec(ctx,
        `INSERT INTO orders (id, user_id, total) VALUES ($1, $2, $3)`,
        orderID, req.UserID, req.Total)
    if err != nil {
        return fmt.Errorf("insert order: %w", err)
    }

    // 2. В той же транзакции записываем событие в Outbox
    payload, _ := json.Marshal(map[string]any{
        "order_id": orderID,
        "user_id":  req.UserID,
        "total":    req.Total,
    })
    _, err = tx.Exec(ctx, `
        INSERT INTO outbox_events (aggregate_type, aggregate_id, event_type, payload)
        VALUES ('Order', $1, 'OrderCreated', $2)
    `, orderID.String(), payload)
    if err != nil {
        return fmt.Errorf("insert outbox: %w", err)
    }

    return tx.Commit(ctx)
}

Интеграция с Redis для буферизации и дедупликации

Redis органично дополняет Outbox Pattern в двух сценариях: буферизация высокочастотных событий и дедупликация на стороне потребителя.

Дедупликация с Redis

Даже при использовании Outbox события могут быть доставлены потребителю более одного раза (at-least-once delivery). На стороне потребителя Redis обеспечивает эффективную дедупликацию через SET NX EX:

func (c *Consumer) isProcessed(ctx context.Context, eventID string) (bool, error) {
    key := "processed_event:" + eventID
    // SET key 1 NX EX 86400 — установить, если не существует, TTL 24 часа
    set, err := c.redis.SetNX(ctx, key, 1, 24*time.Hour).Result()
    if err != nil {
        return false, err
    }
    // set=true означает, что ключ был создан — событие обрабатывается впервые
    return !set, nil
}

func (c *Consumer) Handle(ctx context.Context, event Event) error {
    already, err := c.isProcessed(ctx, event.ID.String())
    if err != nil || already {
        return err // пропускаем дубликат
    }
    return c.processEvent(ctx, event)
}

Redis Streams как промежуточный буфер

При пиковых нагрузках polling worker может не успевать читать из PostgreSQL. Redis Streams (XADD/XREADGROUP) выступают буфером между polling worker и финальным брокером (Kafka, RabbitMQ). Worker публикует в Redis Stream атомарно, а отдельный consumer group пересылает события в Kafka. Это снижает задержку и защищает Kafka от всплесков нагрузки.

Мониторинг и отладка распределённых транзакций

Distributed tracing — обязательный инструмент для систем, использующих Saga и Outbox. Trace ID должен распространяться через все события: записываться в поле metadata таблицы outbox_events и читаться потребителями для восстановления контекста. Используйте OpenTelemetry для инструментирования Go-сервисов.

Ключевые метрики для мониторинга Outbox:

  • outbox_pending_count — количество неопубликованных событий (критический SLO-индикатор);
  • outbox_publish_latency_seconds — задержка между созданием события и его публикацией;
  • outbox_failed_count — события в статусе FAILED, требующие ручного вмешательства;
  • saga_compensation_total — количество выполненных компенсирующих транзакций.

Для диагностики зависших Saga добавьте таблицу saga_state с полем last_updated_at и настройте алёрт при отсутствии обновлений дольше N минут. Это позволит обнаружить «застрявшие» оркестрации до того, как они начнут влиять на пользователей.

PostgreSQL предоставляет мощные инструменты отладки: pg_stat_activity покажет активные блокировки, EXPLAIN ANALYZE — план запросов polling worker. Следите за метрикой autovacuum на таблице outbox_events: при высокой частоте обновлений строк (переход PENDING → PUBLISHED) без своевременного vacuum таблица будет разрастаться.

Сравнение паттернов: когда что применять

Выбор между Saga, Outbox и Event Sourcing не является взаимоисключающим — на практике они часто используются вместе. Однако их цели и trade-off существенно различаются.

  • Saga решает проблему координации бизнес-процессов, охватывающих несколько сервисов. Применяйте Saga везде, где есть многошаговые бизнес-транзакции с возможностью отката. Оркестрацию предпочитайте при сложной логике и необходимости централизованного мониторинга; хореографию — для простых линейных процессов с малым числом участников.
  • Transactional Outbox — универсальный паттерн для надёжной публикации событий. Применяйте его всегда, когда сервис изменяет состояние в базе данных и должен уведомить об этом другие сервисы. Outbox устраняет проблему «двойной записи» и является строительным блоком для любой event-driven архитектуры.
  • Event Sourcing оправдан, когда история изменений является первоклассным бизнес-требованием: финансовые системы, аудит, системы с необходимостью temporal queries. Event Sourcing значительно усложняет систему: сложнее запросы, проблема schema evolution, необходимость снапшотов. Не применяйте его «по умолчанию» только ради модности.

Типичная production-конфигурация для микросервисной системы 2026 года: Saga Orchestration для бизнес-процессов + Transactional Outbox для каждого сервиса + Event Sourcing для агрегатов с богатой историей изменений + Redis для дедупликации и буферизации. PostgreSQL с расширениями pgcrypto и pg_partman (партиционирование таблицы event_store по времени) является надёжной основой для всех трёх паттернов без необходимости введения специализированных хранилищ событий на ранних этапах.

Заключение

Управление состоянием в распределённых системах — одна из фундаментальных проблем backend-архитектуры. Saga, Outbox и Event Sourcing не являются серебряными пулями, но в правильном сочетании они обеспечивают консистентность данных, надёжность доставки событий и полную историю изменений при сохранении независимости сервисов. PostgreSQL в роли надёжного фундамента, Go в роли performant-реализации и Redis в роли быстрого буфера формируют стек, способный обслуживать системы уровня highload с предсказуемым поведением под нагрузкой и прозрачной отладкой.

Технологии

Теги

Руслан Исмаилов

Senior Web / Backend разработчик. Senior web/backend разработчик с 9-летним опытом. Стек: PHP, Laravel, PostgreSQL, Redis, Docker, Kubernetes, REST, микросервисы, CI/CD. Подробнее обо мне →