Архитектурные паттерны Saga, распределённые транзакции, outbox и идемпотентность
0%

Saga, распределённые транзакции, outbox и идемпотентность

Saga, распределённые транзакции, outbox и идемпотентность

В монолите с одной базой оформление заказа выглядит так: открыли транзакцию, записали заказ, уменьшили остаток на складе, списали деньги с внутреннего баланса, закоммитили. Если на любом шаге что-то пошло не так — ROLLBACK, и мира, в котором заказ создан, а деньги не списаны, никогда не существовало. Атомарность нам подарила СУБД, и мы к этому подарку привыкли настолько, что перестали его замечать.

Стоит разнести те же три операции по трём сервисам с тремя базами — и подарок отбирают. Не потому что кто-то поленился реализовать ROLLBACK через сеть, а потому что атомарность между независимыми узлами покупается за доступность, и цена оказывается неприемлемой. Это не инженерная небрежность, это следствие того, как устроены сети и отказы.

Эта статья — про то, что делать вместо потерянной транзакции. Мы разберём:

  • почему 2PC/XA технически работает, но архитектурно почти всегда неверный ответ;
  • что такое сага — цепочка локальных транзакций с семантическими компенсациями;
  • два стиля саг — хореография и оркестрация — и когда какой;
  • аномалии изоляции, которые сага приносит с собой, и контрмеры против них;
  • transactional outbox — как атомарно изменить данные и отправить событие;
  • идемпотентность — почему без неё вся конструкция разваливается, и как её строить честно.

Предполагается, что вы уже знакомы с материалом статей «Микросервисы: границы, коммуникация, данные, эксплуатация» и «Событийная архитектура: брокеры, топики, гарантии доставки» — понятия «at-least-once», «топик», «потребительская группа» здесь используются без повторного введения.


1. Что именно мы потеряли

ACID в одной базе держится на том, что менеджер транзакций видит всё состояние и владеет журналом. Атомарность — это способность одностороннего решения: узел сам решает «коммит» или «откат» и сам его исполняет.

Как только участников несколько, решение становится согласованным: все обязаны прийти к одному исходу. И тут вступает фундаментальная проблема: узел, отправивший «я готов», больше не может решать сам, но и не знает решения. Он заблокирован — держит замки и ждёт.

Формально: атомарный коммит между n узлами — это задача консенсуса, и она сводится к консенсусу в модели с отказами (Gray & Lamport, «Consensus on Transaction Commit», 2004). Классический 2PC — это консенсус с координатором без репликации решения, то есть с единственной точкой отказа.

Есть и более приземлённая арифметика. Пусть каждый сервис доступен с вероятностью 0.999. Транзакция, требующая одновременной доступности пяти сервисов, доступна с вероятностью 0.999⁵ ≈ 0.995 — то есть 0.5% отказов вместо 0.1%. Синхронная атомарность перемножает доступности; асинхронная сага — нет, потому что временно недоступный участник просто получит сообщение позже.

1.1. Как работает 2PC и где он ломается

Ключевые следствия:

  1. 2PC — блокирующий протокол. Отказ координатора между фазами оставляет участников in-doubt. Никакой таймаут не спасает: откатиться нельзя (вдруг решение было COMMIT), закоммитить нельзя (вдруг ROLLBACK). Существует 3PC, снимающий блокировку ценой предположения о синхронной сети, — на практике не применяется.
  2. Длительность удержания замков равна сетевому RTT × 2 плюс время самого медленного участника. В однобазовой транзакции замок живёт микросекунды; в XA — десятки миллисекунд и больше. Пропускная способность по «горячим» строкам падает на порядок.
  3. Требуется поддержка XA у всех участников. Kafka, большинство NoSQL-хранилищ, внешние платёжные API и HTTP-сервисы её не имеют. XA живёт в мире «Java EE + пара реляционных СУБД + JMS».
  4. Он расползается по границам владения. Чтобы участвовать в чужой транзакции, сервис должен пустить чужой координатор в свой жизненный цикл замков — это прямое нарушение автономии, ради которой сервисы и разделяли.

Пэт Хелланд сформулировал это ещё в 2007 году в статье «Life beyond Distributed Transactions: An Apostate’s Opinion»: в масштабируемых системах программист обязан работать с «почти-транзакциями» поверх сущностей, помещающихся в один узел, и явным образом обрабатывать неопределённость между ними. Это программа-минимум всей оставшейся статьи.

Когда 2PC всё же уместен: несколько схем внутри одного кластера СУБД; интеграция с легаси-мейнфреймом, где иначе никак; Kafka-транзакции внутри одного кластера (read-process-write). То есть — когда участники технологически однородны, находятся рядом и отказ координатора восстанавливается за секунды. Между бизнес-сервисами в разных командах — почти никогда.


2. Сага: обмен отката на компенсацию

Идея саги старше микросервисов на три десятилетия. В 1987 году Гектор Гарсиа-Молина и Кеннет Салем опубликовали работу «Sagas» (SIGMOD'87). Задача была другая — длительные транзакции (LLT), держащие замки часами и мешающие всем остальным. Решение оказалось универсальным:

Сага — последовательность локальных транзакций T₁…Tₙ, каждая из которых коммитится независимо, снабжённая компенсирующими транзакциями C₁…Cₙ₋₁. Если Tₖ провалилась, система выполняет Cₖ₋₁, Cₖ₋₂, …, C₁, приводя систему в семантически согласованное состояние.

Сага гарантирует ACD — атомарность (в смысле «всё или семантически ничего»), согласованность, долговечность. Изоляции (I) у саги нет — и это главное, что нужно про неё понять. Промежуточные состояния саги видимы всем.

2.1. Компенсация — не откат

Это различие стоит вбить гвоздями, потому что 80% ошибок в сагах растут отсюда.

ROLLBACK Компенсация
Уровень физический (журнал СУБД) семантический (доменное действие)
Видимость эффекта никто никогда не видел все видели и могли отреагировать
Результат состояния не было появилось новое состояние, отменяющее прежнее
Пример строка не изменилась «списание 1000 ₽» + «возврат 1000 ₽» — две записи в выписке
Может ли провалиться нет да, и это надо проектировать

Классический пример неотменяемости: письмо клиенту отправлено. Компенсации «не отправлять письмо» не существует — существует только «отправить письмо с извинением». Поэтому шаги вроде отправки уведомлений выносят в самый конец саги, после точки невозврата.

2.2. Три класса шагов и точка невозврата

Крис Ричардсон в «Microservices Patterns» вводит классификацию, без которой сагу невозможно спроектировать корректно:

  • compensatable — шаг, для которого написана компенсация; выполняется до pivot;
  • pivot — точка невозврата: если он прошёл, сага обязана дойти до конца; если провалился — откатываем всё, что было до;
  • retriable — шаг после pivot; он не имеет права провалиться навсегда, только «пока не получилось, повторим».

Три класса шагов саги: компенсируемые, pivot, повторяемые

Практическое правило проектирования: расположите шаги так, чтобы все опасные и необратимые операции оказались как можно правее, а pivot был ровно один. Если у вас два необратимых шага подряд (списать деньги в одном провайдере и списать бонусы в другом) — вы обязаны либо сделать один из них компенсируемым (возврат), либо объединить их в один сервис с одной локальной транзакцией.

2.3. Жизненный цикл саги

Обратите внимание на состояние STUCK. Его почти всегда забывают, и это худшая из ошибок: retriable-шаг после pivot не может быть отменён, значит, при исчерпании повторов сага не имеет права молча умереть. Она обязана попасть в очередь на разбор человеком, с алертом. Деньги списаны — отгрузки нет; это инцидент, а не «сообщение в DLQ».


3. Хореография против оркестрации

Сагу можно координировать двумя способами.

3.1. Хореография: координатора нет

Каждый сервис слушает события других и публикует свои. Логика саги «размазана» по подписчикам.

Плюсы: нет единой точки отказа и единой точки изменений; сервисы связаны только контрактами событий; добавить нового участника — значит просто подписаться.

Минусы, растущие нелинейно:

  • Логики саги не существует ни в одном файле. Чтобы ответить на вопрос «что происходит при отказе оплаты», нужно прочитать три сервиса и построить граф в голове.
  • Циклические зависимости событий. Order слушает Payment, Payment слушает Inventory, Inventory слушает Order — и вот у вас распределённый цикл, который никто не видит.
  • Нет места, где хранится состояние саги целиком. «Сколько саг сейчас висит на шаге оплаты?» — вопрос без ответа.
  • Компенсации становятся отдельным набором событий, и их корректность никто не проверяет.

Хореография хороша для саг из 2–3 шагов и там, где шаги действительно независимы. На 5+ шагах она превращается в то, что Сэм Ньюмен называет «эмерджентным поведением, которое никто не проектировал».

3.2. Оркестрация: явная машина состояний

Появляется оркестратор — компонент (обычно живущий в сервисе-инициаторе), который хранит состояние саги и явно рассылает команды.

Плюсы: логика в одном месте, читается как код; состояние саги персистентно и наблюдаемо; таймауты и повторы централизованы; тестируется как обычный конечный автомат.

Минус: оркестратор — новый компонент, который легко превращается в «божественный сервис», знающий бизнес-логику чужих доменов. Лечится дисциплиной: оркестратор знает последовательность и условия перехода, но не знает, как считается НДС.

Практическое правило: до 3 шагов и при отсутствии компенсаций — хореография; от 4 шагов или при наличии компенсаций — оркестрация. Ричардсон в microservices.io/patterns/data/saga даёт ту же рекомендацию.


4. Реализация оркестратора

Сначала — псевдокод исполнения саги. Он короткий, и в нём вся суть.

исполнить_сагу(saga):
    пока saga.step < len(steps):
        шаг = steps[saga.step]
        если saga.direction == FORWARD:
            результат = вызвать(шаг.action, saga.context, ключ=idem_key(saga, шаг))
            если результат.ok:
                saga.context.обновить(результат.data)
                saga.step += 1
            иначе если шаг.retriable:
                запланировать_повтор(saga, backoff)      # после pivot — только вперёд
                вернуть
            иначе:
                saga.direction = BACKWARD
                saga.step -= 1
        иначе:                                            # BACKWARD
            если шаг.compensation is None:
                saga.step -= 1                            # шаг не требует компенсации
            иначе:
                результат = вызвать(шаг.compensation, saga.context, ключ=comp_key(saga, шаг))
                если результат.ok:
                    saga.step -= 1
                иначе:
                    запланировать_повтор(saga, backoff)   # компенсация ОБЯЗАНА пройти
                    вернуть
        сохранить(saga)                                   # персист после КАЖДОГО перехода
    saga.state = COMPLETED если saga.direction == FORWARD иначе CANCELLED

Два инварианта, которые нельзя нарушать:

  1. Состояние саги сохраняется после каждого перехода, до отправки следующей команды. Иначе после падения оркестратор не знает, где он был.
  2. Компенсации не имеют права сдаться. Они повторяются бесконечно с backoff; при исчерпании разумного лимита — алерт человеку, но не «отмена отмены».

Теперь рабочая реализация на Python. Шаги описываются декларативно, состояние живёт в БД, отправка команд идёт через outbox (о нём — в разделе 6).

from __future__ import annotations
import json
import uuid
from dataclasses import dataclass, field
from enum import Enum
from typing import Callable, Optional

class Direction(str, Enum):
    FORWARD = "FORWARD"
    BACKWARD = "BACKWARD"

@dataclass(frozen=True)
class Step:
    name: str
    # действие возвращает команду для отправки участнику (или None, если шаг локальный)
    action: Callable[[dict], dict]
    compensation: Optional[Callable[[dict], dict]] = None
    # pivot: после него откат невозможен, только повторы вперёд
    pivot: bool = False
    max_attempts: int = 12

@dataclass
class SagaInstance:
    saga_id: uuid.UUID
    saga_type: str
    step_index: int = 0
    direction: Direction = Direction.FORWARD
    state: str = "RUNNING"
    attempts: int = 0
    context: dict = field(default_factory=dict)

class SagaDefinition:
    def __init__(self, name: str, steps: list[Step]):
        pivots = [i for i, s in enumerate(steps) if s.pivot]
        if len(pivots) > 1:
            raise ValueError("В саге допустим ровно один pivot-шаг")
        # шаги ПОСЛЕ pivot не должны иметь компенсаций — их некуда откатывать
        if pivots:
            for s in steps[pivots[0] + 1:]:
                if s.compensation is not None:
                    raise ValueError(f"Шаг {s.name} после pivot не может быть компенсируемым")
        self.name, self.steps = name, steps

class SagaEngine:
    """Продвигает сагу ровно на один переход за вызов.

    Вызывается двумя источниками: (1) ответом участника, (2) планировщиком повторов.
    Вся работа идёт внутри ОДНОЙ локальной транзакции сервиса-оркестратора:
    изменение saga_instance и запись команды в outbox атомарны.
    """

    def __init__(self, definition: SagaDefinition, repo, outbox):
        self.d, self.repo, self.outbox = definition, repo, outbox

    def handle_reply(self, saga_id: uuid.UUID, ok: bool, payload: dict) -> None:
        with self.repo.transaction() as tx:
            # SELECT ... FOR UPDATE: сага — критическая секция, параллельные ответы недопустимы
            saga = self.repo.load_for_update(tx, saga_id)
            if saga.state != "RUNNING":
                return                                   # поздний дубликат ответа — игнорируем
            step = self.d.steps[saga.step_index]

            if saga.direction is Direction.FORWARD:
                if ok:
                    saga.context.update(payload)
                    saga.step_index += 1
                    saga.attempts = 0
                elif step.pivot or self._after_pivot(saga.step_index):
                    self._schedule_retry(tx, saga, step)  # назад нельзя
                    return
                else:
                    saga.direction = Direction.BACKWARD
                    saga.step_index -= 1
                    saga.attempts = 0
            else:                                         # BACKWARD
                if ok:
                    saga.step_index -= 1
                    saga.attempts = 0
                else:
                    self._schedule_retry(tx, saga, step)  # компенсация обязана пройти
                    return

            self._advance(tx, saga)

    def _advance(self, tx, saga: SagaInstance) -> None:
        if saga.direction is Direction.FORWARD and saga.step_index >= len(self.d.steps):
            saga.state = "COMPLETED"
        elif saga.direction is Direction.BACKWARD and saga.step_index < 0:
            saga.state = "CANCELLED"
        else:
            step = self.d.steps[saga.step_index]
            fn = step.action if saga.direction is Direction.FORWARD else step.compensation
            if fn is None:                                 # нечего компенсировать — идём дальше назад
                saga.step_index -= 1
                return self._advance(tx, saga)
            command = fn(saga.context)
            # ключ идемпотентности детерминирован: (сага, шаг, направление).
            # Повтор того же шага после падения даст ТОТ ЖЕ ключ — участник распознает дубль.
            command["idempotency_key"] = f"{saga.saga_id}:{saga.step_index}:{saga.direction.value}"
            self.outbox.enqueue(tx, topic=command.pop("topic"), payload=command)
        self.repo.save(tx, saga)

    def _after_pivot(self, idx: int) -> bool:
        return any(s.pivot for s in self.d.steps[:idx])

    def _schedule_retry(self, tx, saga: SagaInstance, step: Step) -> None:
        saga.attempts += 1
        if saga.attempts >= step.max_attempts:
            saga.state = "STUCK"                          # алерт человеку, сага не умирает молча
            self.repo.save(tx, saga)
            self.outbox.enqueue(tx, topic="saga.stuck",
                                payload={"saga_id": str(saga.saga_id), "step": step.name})
            return
        delay = min(2 ** saga.attempts, 3600)             # экспоненциальный backoff, потолок 1 час
        self.repo.schedule(tx, saga.saga_id, delay_seconds=delay)
        self.repo.save(tx, saga)

Определение конкретной саги получается декларативным и читаемым:

order_saga = SagaDefinition("create_order", [
    Step("reserve_stock",
         action=lambda c: {"topic": "inventory.commands",
                           "type": "ReserveStock", "order_id": c["order_id"], "items": c["items"]},
         compensation=lambda c: {"topic": "inventory.commands",
                                 "type": "ReleaseStock", "order_id": c["order_id"]}),
    Step("charge_card",
         action=lambda c: {"topic": "payment.commands",
                           "type": "ChargeCard", "order_id": c["order_id"], "amount": c["total"]},
         pivot=True),                                   # точка невозврата
    Step("create_shipment",
         action=lambda c: {"topic": "shipping.commands",
                           "type": "CreateShipment", "order_id": c["order_id"]}),
    Step("notify_customer",
         action=lambda c: {"topic": "notification.commands",
                           "type": "SendOrderConfirmation", "order_id": c["order_id"]}),
])

Сложность. Пусть n — число шагов саги. Успешный проход: n локальных транзакций, 2n сообщений (команда + ответ), O(n) записей состояния. Худший случай (отказ на предпоследнем компенсируемом шаге): n + (n−1) локальных транзакций и до 4n сообщений. Память на сагу — O(размер контекста); контекст обязан оставаться маленьким (идентификаторы и суммы, не корзины целиком), иначе таблица саг становится вторым хранилищем данных. По латентности сага — сумма шагов, а не максимум: сага из 5 шагов по 50 мс — это 250 мс, и это ещё одна причина держать n маленьким.


5. Изоляции нет: аномалии и контрмеры

Самая недооценённая часть темы. Сага коммитит каждый шаг сразу, поэтому промежуточные состояния видны всем. Возникают ровно те аномалии, от которых нас защищал уровень изоляции СУБД:

  • Грязное чтение (dirty read). Другой процесс видит заказ в состоянии PENDING со снятым резервом — и принимает решение на основе состояния, которое через секунду будет отменено. Отчёт «продажи за час» посчитает заказ, который вот-вот отменится.
  • Потерянные обновления (lost update). Сага читает баланс, считает, пишет. Параллельно пользователь тратит деньги через другой путь. Компенсация «вернуть 1000 ₽» пишет абсолютное значение и затирает чужую запись.
  • Неповторяющееся чтение (fuzzy read). Сага прочитала цену на шаге 1, применила на шаге 4 — за это время цена изменилась.

Контрмеры (терминология из «Microservices Patterns», глава 4):

Контрмера Суть Когда применять
Semantic lock Флаг «в процессе»: order.status = PENDING, account.pending_debit. Читатели обязаны его учитывать Базовая мера, нужна почти всегда
Commutative updates Компенсация коммутативна с прочими операциями: balance = balance + 1000 вместо balance = 5000 Всегда, где возможно — снимает lost update
Pessimistic view Переупорядочить шаги так, чтобы рискованное действие было позже: не начислять бонусы до подтверждения оплаты Когда цена грязного чтения высока
Re-read value Перед записью перечитать и сверить версию (optimistic offline lock) Обновление изменчивых сущностей
Version file Записывать все операции над сущностью с версиями и уметь применять их не по порядку Когда команды приходят вне порядка
By value Выбирать механизм по риску: мелкие суммы — сага, крупные переводы — 2PC/ручное подтверждение Финансы, регуляторика

Иллюстрация semantic lock и коммутативного обновления на SQL:

-- ПЛОХО: абсолютная запись. Компенсация затрёт параллельное изменение.
UPDATE accounts SET balance = 5000 WHERE id = 42;

-- ХОРОШО: коммутативная дельта. Порядок применения не важен,
-- параллельные операции складываются корректно.
UPDATE accounts SET balance = balance + 1000 WHERE id = 42;

-- Semantic lock: резерв средств отдельным полем.
-- Доступный баланс = balance - held. Все читатели обязаны считать так.
UPDATE accounts
   SET held = held + :amount
 WHERE id = :account_id
   AND balance - held >= :amount;     -- защита от овердрафта в одном атомарном шаге
-- 0 затронутых строк = недостаточно средств, шаг саги провалился

-- Компенсация: снять резерв. Идемпотентна за счёт проверки в hold-таблице.
UPDATE accounts SET held = held - :amount WHERE id = :account_id;

Отдельно про семантические блокировки и живучесть: любой PENDING обязан иметь таймаут. Сага, застрявшая на шаге оплаты, держит товар в резерве; без TTL склад через неделю окажется полностью «зарезервирован» призраками. Реализуется отдельным процессом-сборщиком, который ищет саги старше N минут и запускает компенсацию (это тоже переход в машине состояний, а не отдельный скрипт «почистить базу»).


6. Transactional outbox: проблема двойной записи

Вернёмся к оркестратору. В коде выше запись состояния саги и постановка команды шли в одной транзакции — через outbox.enqueue(tx, ...). Это не деталь стиля, это единственный способ не потерять шаг.

Проблема формулируется так: сервису нужно атомарно изменить своё состояние в БД и отправить сообщение в брокер. Это две разные системы, общей транзакции нет.

Проблема двойной записи: два порядка, обе дырявые

Обходные пути, которые не работают:

  • «Отправим после коммита, в try/finally» — процесс может умереть между коммитом и отправкой. Окно маленькое, но при 10 000 RPS «маленькое» означает несколько потерь в день.
  • «Обернём в XA» — Kafka не поддерживает XA-протокол с внешним координатором, и мы возвращаемся к разделу 1.
  • «Публикуем внутри транзакции, при ошибке откатываем» — сообщение уже ушло, откатить его нельзя.

6.1. Решение: событие как строка таблицы

Записываем сообщение в таблицу outbox той же транзакцией, что и бизнес-данные. Отдельный процесс (relay) читает таблицу и публикует в брокер.

CREATE TABLE outbox_message (
    id             BIGSERIAL PRIMARY KEY,          -- монотонный порядок публикации
    message_id     UUID        NOT NULL UNIQUE,    -- то, по чему дедуплицирует получатель
    aggregate_type TEXT        NOT NULL,
    aggregate_id   TEXT        NOT NULL,           -- станет ключом партиции в Kafka
    topic          TEXT        NOT NULL,
    payload        JSONB       NOT NULL,
    created_at     TIMESTAMPTZ NOT NULL DEFAULT now(),
    published_at   TIMESTAMPTZ                     -- NULL = ждёт отправки
);

-- Частичный индекс: в него попадают ТОЛЬКО неотправленные строки.
-- Таблица может быть огромной, индекс останется крошечным.
CREATE INDEX outbox_pending_idx ON outbox_message (id) WHERE published_at IS NULL;

Бизнес-операция и событие пишутся вместе:

def approve_order(conn, order_id: uuid.UUID) -> None:
    with conn.transaction():                      # ОДНА локальная ACID-транзакция
        conn.execute(
            "UPDATE orders SET status = 'APPROVED', updated_at = now() WHERE id = %s",
            (order_id,),
        )
        conn.execute(
            """INSERT INTO outbox_message
                   (message_id, aggregate_type, aggregate_id, topic, payload)
               VALUES (%s, 'Order', %s, 'order.events', %s)""",
            (uuid.uuid4(), str(order_id),
             json.dumps({"type": "OrderApproved", "order_id": str(order_id)})),
        )
    # Коммит атомарен: либо есть и изменение, и событие, либо нет ни того, ни другого.

6.2. Relay: как вынимать и публиковать

Два способа.

А. Polling publisher. Простой цикл с SKIP LOCKED — работает на любой СУБД, запускается в нескольких экземплярах.

BATCH = 200

def relay_once(conn, producer) -> int:
    with conn.transaction():
        rows = conn.execute(
            """SELECT id, message_id, topic, aggregate_id, payload
                 FROM outbox_message
                WHERE published_at IS NULL
                ORDER BY id                       -- строгий порядок = порядок коммитов
                LIMIT %s
                FOR UPDATE SKIP LOCKED""",        -- параллельные relay не мешают друг другу
            (BATCH,),
        ).fetchall()

        for r in rows:
            # Ключ = aggregate_id: все события одного агрегата попадут в одну партицию
            # и придут потребителю в правильном порядке.
            producer.send(
                topic=r["topic"],
                key=r["aggregate_id"].encode(),
                value=json.dumps(r["payload"]).encode(),
                headers=[("message_id", str(r["message_id"]).encode())],
            )
        producer.flush()                          # подтверждение брокера ДО коммита пометки

        if rows:
            conn.execute("UPDATE outbox_message SET published_at = now() WHERE id = ANY(%s)",
                         ([r["id"] for r in rows],))
    return len(rows)

Порядок здесь принципиален: сначала flush() (брокер подтвердил приём), потом пометка published_at. Если процесс упадёт между ними — сообщение уйдёт повторно. Это осознанный выбор в пользу at-least-once: дубликат мы вылечим идемпотентностью, потерю — ничем.

Плюсы: тривиально, никакой инфраструктуры. Минусы: задержка = интервал опроса; нагрузка на БД растёт с частотой; при высоком трафике нужен разумный batch. SKIP LOCKED даёт горизонтальное масштабирование, но ломает глобальный порядок между воркерами — поэтому порядок гарантируется только внутри партиции по aggregate_id, что почти всегда и требуется.

Б. Transaction log tailing (CDC). Debezium читает WAL/binlog СУБД и публикует изменения. Специальный SMT Outbox Event Router превращает строки outbox-таблицы в нормальные события нужного топика:

# Коннектор Debezium: только outbox-таблица, маршрутизация по полю topic
name: order-outbox-connector
config:
  connector.class: io.debezium.connector.postgresql.PostgresConnector
  plugin.name: pgoutput
  table.include.list: public.outbox_message
  tombstones.on.delete: "false"
  transforms: outbox
  transforms.outbox.type: io.debezium.transforms.outbox.EventRouter
  transforms.outbox.table.field.event.id: message_id      # → header id для дедупликации
  transforms.outbox.table.field.event.key: aggregate_id   # → ключ партиции
  transforms.outbox.table.field.event.payload: payload
  transforms.outbox.route.by.field: topic                 # маршрут берём из строки

Плюсы CDC: задержка в миллисекундах, нулевая дополнительная нагрузка на таблицы (читается журнал, который и так пишется), строгий порядок коммитов. Минусы: Kafka Connect в эксплуатации, чувствительность к миграциям схемы, необходимость следить за слотом репликации (забытый слот в PostgreSQL не даёт удалять WAL и однажды заполняет диск — классическая ночная авария).

Практика: начинать с polling publisher (100–200 мс интервала хватает почти всем), переезжать на CDC, когда задержка или нагрузка станут проблемой. Оба варианта требуют джоба-уборщика, удаляющего строки с published_at < now() - interval '7 days' — иначе таблица растёт вечно.

6.3. Обратная сторона: inbox

Симметричная проблема на стороне получателя: сообщение обработано, но commit offset не прошёл — брокер пришлёт его снова. Отсюда — inbox / idempotent consumer: получатель в той же транзакции, что и бизнес-эффект, записывает message_id в таблицу с UNIQUE-ограничением.

def handle_message(conn, message_id: uuid.UUID, group: str, payload: dict) -> None:
    try:
        with conn.transaction():
            # Ставка на ограничение БД, а не на SELECT-then-INSERT:
            # проверка "видел ли я это раньше" и запись эффекта атомарны.
            conn.execute(
                "INSERT INTO inbox_message (message_id, consumer_group) VALUES (%s, %s)",
                (message_id, group),
            )
            apply_business_effect(conn, payload)      # тот же tx — либо всё, либо ничего
    except UniqueViolation:
        # Дубликат: эффект уже применён в прошлый раз. Тихо подтверждаем сообщение.
        log.debug("duplicate message %s ignored", message_id)

Почему не SELECT ... WHERE message_id = ? перед вставкой: два конкурентных потребителя пройдут проверку одновременно и оба применят эффект. Уникальный индекс — единственная проверка, атомарная относительно параллелизма.

Итоговая цепочка outbox → брокер → inbox даёт то, что называют effectively-once: физически сообщение доставляется at-least-once, а наблюдаемый эффект применяется ровно один раз.


7. Идемпотентность

Идемпотентная операция — та, повторное применение которой не меняет результат: f(f(x)) = f(x). Без неё все повторы, retry и at-least-once превращаются из механизма надёжности в механизм порчи данных.

7.1. Что идемпотентно само по себе, а что нет

Операция Идемпотентна? Почему
SET status = 'PAID' да абсолютное присваивание
balance = balance - 100 нет накопительный эффект
INSERT с натуральным PK да второй раз конфликт ключа
INSERT с автоинкрементом нет появится вторая строка
DELETE WHERE id = 42 да второй раз 0 строк
HTTP PUT /orders/42 да (по контракту) полная замена ресурса
HTTP POST /charges нет создаёт новый ресурс каждый раз
Отправка письма нет внешний необратимый эффект

Важная тонкость: идемпотентность ≠ коммутативность. Операция может быть идемпотентной, но чувствительной к порядку: SET status='SHIPPED' и SET status='CANCELLED' обе идемпотентны, но результат зависит от того, какая пришла последней. Если брокер не гарантирует порядок между партициями, нужен ещё и fencing — монотонный номер, позволяющий отбросить устаревшую команду:

-- Обновление применяется, только если версия события новее уже применённой.
UPDATE orders
   SET status = :new_status,
       last_event_version = :event_version
 WHERE id = :order_id
   AND last_event_version < :event_version;   -- старое событие просто не сработает

Мартин Клеппман подробно разбирает необходимость fencing-токенов в «How to do distributed locking» — там речь о блокировках, но механика защиты от «отставшего» клиента ровно та же.

7.2. Ключи идемпотентности в API

Когда операция принципиально неидемпотентна (создать платёж), её делают идемпотентной снаружи: клиент передаёт ключ, сервер запоминает результат первого выполнения и возвращает его на повторы. Так устроен Stripe с заголовком Idempotency-Key; сейчас механизм стандартизуется в IETF-черновике Idempotency-Key header field.

CREATE TABLE idempotency_key (
    key            TEXT PRIMARY KEY,
    request_hash   TEXT        NOT NULL,   -- отпечаток тела запроса
    state          TEXT        NOT NULL,   -- IN_PROGRESS | COMPLETED
    response_code  INT,
    response_body  JSONB,
    created_at     TIMESTAMPTZ NOT NULL DEFAULT now(),
    expires_at     TIMESTAMPTZ NOT NULL    -- Stripe хранит 24 часа
);
def charge(conn, idem_key: str, request: dict):
    body_hash = sha256(canonical_json(request))

    try:
        with conn.transaction():
            conn.execute(
                """INSERT INTO idempotency_key (key, request_hash, state, expires_at)
                   VALUES (%s, %s, 'IN_PROGRESS', now() + interval '24 hours')""",
                (idem_key, body_hash),
            )
            # Захватили ключ первыми — выполняем и фиксируем результат в ТОЙ ЖЕ транзакции.
            result = do_charge(conn, request)
            conn.execute(
                """UPDATE idempotency_key
                      SET state = 'COMPLETED', response_code = 201, response_body = %s
                    WHERE key = %s""",
                (json.dumps(result), idem_key),
            )
            return 201, result

    except UniqueViolation:
        row = conn.execute("SELECT * FROM idempotency_key WHERE key = %s", (idem_key,)).fetchone()

        # Тот же ключ с ДРУГИМ телом — почти всегда баг клиента. Молчать нельзя.
        if row["request_hash"] != body_hash:
            return 422, {"error": "idempotency_key_reuse_with_different_payload"}

        if row["state"] == "IN_PROGRESS":
            # Первый запрос ещё выполняется. Просим клиента повторить позже,
            # а не выполняем параллельно.
            return 409, {"error": "request_in_progress", "retry_after": 1}

        return row["response_code"], row["response_body"]

Три детали, которые обычно забывают и которые отличают рабочую реализацию от игрушечной:

  1. Хэш тела запроса. Иначе клиент, переиспользовавший ключ для другой суммы, получит ответ от старого платежа и будет уверен, что новый прошёл.
  2. Состояние IN_PROGRESS. Без него два одновременных повтора выполнят операцию дважды: оба не найдут завершённой записи.
  3. TTL. Хранить ключи вечно нельзя (таблица растёт), хранить слишком мало — опасно (клиент повторит после долгой сетевой паузы). Сутки — разумная отраслевая норма.

Кто генерирует ключ: клиент, один раз на логическую операцию, и переиспользует его на всех повторах. Ключ, сгенерированный на каждую попытку заново, бесполезен. В саге ключ детерминирован — {saga_id}:{step}:{direction} (см. код в разделе 4), поэтому повтор шага после падения оркестратора автоматически распознаётся как дубль.

7.3. «Exactly-once» и что за этим стоит

Строгая доставка ровно-один-раз в асинхронной сети с отказами невозможна: это следствие проблемы двух генералов. Отправитель, не получивший подтверждения, не может отличить «сообщение не дошло» от «дошло, а подтверждение потерялось», и обязан выбрать: повторить (риск дубля) или нет (риск потери).

Kafka-транзакции с enable.idempotence=true и isolation.level=read_committed дают exactly-once processing в контуре read-process-write: атомарно фиксируются и результат обработки, и смещения потребителя (документация Kafka, детали — в разборе Confluent). Но это работает только внутри Kafka. Как только побочный эффект уходит наружу — в вашу БД, в платёжный шлюз, в SMTP — гарантия заканчивается, и остаётся единственная формула:

effectively-once = at-least-once доставка + идемпотентный получатель.

Всё, что описано в этой статье, — способы реализовать вторую половину этой формулы.


8. Как это выглядит в проде

Готовые движки саг. Писать оркестратор руками стоит, только если саг мало и они простые. Иначе:

  • Temporal (и его предок Cadence из Uber) — саги пишутся как обычный императивный код на Go/Java/TypeScript/Python; состояние восстанавливается детерминированным воспроизведением истории событий. Компенсации выражаются через defer/try-finally и читаются как в однопроцессной программе.
  • AWS Step Functions — машина состояний в JSON, с Catch/Retry, встроенными таймаутами и визуализацией. Хорошо ложится на serverless (см. «Serverless и edge-архитектуры»).
  • Camunda 8 / Zeebe — BPMN-модель, когда процесс должен обсуждаться с бизнесом и быть видимым нетехническим людям.
  • Eventuate Tram Sagas, MassTransit Saga State Machine, NServiceBus Sagas — библиотечные фреймворки для Java/.NET.
  • MicroProfile LRA — стандартизованные «долгоживущие действия» с аннотациями @LRA/@Compensate.

Наблюдаемость. Сага живёт минутами и часами, обычный трейс её не покрывает. Что нужно обязательно:

  • сквозной saga_id, проставляемый в каждое сообщение и каждый лог-запись (пробрасывается вместе с trace context);
  • метрики: saga_duration_seconds по типу саги, счётчики по терминальным состояниям (completed/cancelled/stuck), возраст самой старой RUNNING-саги — лучший индикатор «что-то зависло»;
  • алерт на state = 'STUCK' и на рост outbox (неотправленные строки старше минуты = relay сломан);
  • админ-интерфейс для просмотра и ручного продвижения саги: инцидент «деньги списаны, отгрузки нет» решается человеком, и ему нужен инструмент.

Типичные ошибки (по частоте, с которой встречаются в ревью):

  1. Компенсация не идемпотентна. ReleaseStock вызывается дважды и возвращает товар на склад дважды. Лечится: ключ идемпотентности + inbox на стороне участника.
  2. Нет состояния STUCK. Сага молча уходит в DLQ после pivot, деньги списаны, никто не узнал.
  3. Два pivot-шага. Два необратимых внешних вызова — гарантированная ручная работа при отказе второго. Проверяйте на этапе описания саги (в коде выше это делает SagaDefinition.__init__).
  4. Отсутствие TTL у semantic lock. Резервы накапливаются, склад «кончается» при полных полках.
  5. Толстый контекст саги. В context кладут всю корзину, документы, адреса. Таблица саг превращается во второе хранилище, миграции становятся адом. Кладите идентификаторы.
  6. Публикация в брокер вместо outbox «потому что и так работает». Работает, пока не упадёт под нагрузкой — тогда потери обнаружатся через неделю по расхождению отчётов.
  7. Оркестратор знает чужую логику. Считает скидки, валидирует адреса. Возвращаемся к распределённому монолиту — см. «Монолит и модульный монолит».
  8. Компенсация пишет абсолютное значение вместо дельты — lost update при параллельных операциях.

9. Как выбрать механизм

Обратите внимание на две зелёные ветки вверху. Лучшая распределённая транзакция — та, которой нет. Если операция постоянно требует согласованного изменения в трёх сервисах, это чаще всего сигнал о неверно проведённой границе: три сервиса на самом деле являются одним ограниченным контекстом. Прежде чем строить сагу, всерьёз рассмотрите перенос границы — материал о том, как их искать, есть в статьях «Микросервисы» и в треке по DDD.


Мини-итог

  • Атомарность между узлами покупается за доступность. 2PC блокирующий, перемножает доступности участников и нарушает автономию сервисов; между бизнес-сервисами он почти всегда неверный ответ.
  • Сага — цепочка локальных транзакций с семантическими компенсациями. Даёт ACD; изоляции нет, промежуточные состояния видны всем.
  • Компенсация ≠ откат: она видима, доменна и сама может провалиться. Проектируйте шаги как compensatable → pivot → retriable, с ровно одним pivot.
  • Хореография — до 3 простых шагов; оркестрация — от 4 шагов или при наличии компенсаций. Состояние саги персистентно и наблюдаемо, STUCK — обязательное состояние.
  • Аномалии изоляции лечатся контрмерами: semantic lock (с обязательным TTL), коммутативные обновления, перестановка шагов, перечитывание с версией.
  • Transactional outbox решает двойную запись: бизнес-данные и событие пишутся одной локальной транзакцией, relay (polling или CDC) публикует их в брокер.
  • Идемпотентность — фундамент всей конструкции: inbox с UNIQUE-ограничением у потребителя, ключи идемпотентности с хэшем тела и состоянием IN_PROGRESS в API, fencing по версии против нарушения порядка.
  • Exactly-once снаружи одной системы не существует. Работает только формула: at-least-once доставка + идемпотентный получатель.

Источники


Что дальше

Сага и outbox описывают, как сервисы согласуют изменения между собой. Следующий вопрос — через какой интерфейс они вообще разговаривают: синхронный REST, типизированный gRPC, гибкий GraphQL или вебхуки наружу, и как менять эти контракты, не ломая потребителей.

Читайте дальше: «Стили API: REST, GraphQL, gRPC, вебхуки и версионирование».

Нашли неточность? Выделите фрагмент текста — рядом появится жучок.

Нужен разбор именно вашей ситуации?

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

Доска запросов