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 и где он ломается
С этого момента решение необратимо. Note over C,S: Фаза 2 — исполнение C->>O: commit(tx-42) C--xP: commit(tx-42) — координатор упал Note over P,S: Payment и Inventory держат замки
и НЕ МОГУТ ни закоммитить, ни откатиться:
они не знают решения. Это in-doubt состояние. Note over P,S: Замки висят до восстановления координатора.
Все, кто ждёт эти строки, стоят в очереди.
Ключевые следствия:
- 2PC — блокирующий протокол. Отказ координатора между фазами оставляет участников in-doubt. Никакой таймаут не спасает: откатиться нельзя (вдруг решение было COMMIT), закоммитить нельзя (вдруг ROLLBACK). Существует 3PC, снимающий блокировку ценой предположения о синхронной сети, — на практике не применяется.
- Длительность удержания замков равна сетевому RTT × 2 плюс время самого медленного участника. В однобазовой транзакции замок живёт микросекунды; в XA — десятки миллисекунд и больше. Пропускная способность по «горячим» строкам падает на порядок.
- Требуется поддержка XA у всех участников. Kafka, большинство NoSQL-хранилищ, внешние платёжные API и HTTP-сервисы её не имеют. XA живёт в мире «Java EE + пара реляционных СУБД + JMS».
- Он расползается по границам владения. Чтобы участвовать в чужой транзакции, сервис должен пустить чужой координатор в свой жизненный цикл замков — это прямое нарушение автономии, ради которой сервисы и разделяли.
Пэт Хелланд сформулировал это ещё в 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 был ровно один. Если у вас два необратимых шага подряд (списать деньги в одном провайдере и списать бонусы в другом) — вы обязаны либо сделать один из них компенсируемым (возврат), либо объединить их в один сервис с одной локальной транзакцией.
2.3. Жизненный цикл саги
Обратите внимание на состояние STUCK. Его почти всегда забывают, и это худшая из ошибок: retriable-шаг после pivot не может быть отменён, значит, при исчерпании повторов сага не имеет права молча умереть. Она обязана попасть в очередь на разбор человеком, с алертом. Деньги списаны — отгрузки нет; это инцидент, а не «сообщение в DLQ».
3. Хореография против оркестрации
Сагу можно координировать двумя способами.
3.1. Хореография: координатора нет
Каждый сервис слушает события других и публикует свои. Логика саги «размазана» по подписчикам.
(PENDING)"] A2["OrderApproved"] A3["OrderCancelled"] end subgraph IS["Inventory Service"] B["StockReserved"] B2["StockRejected"] B3["StockReleased"] end subgraph PS["Payment Service"] C["PaymentCharged"] C2["PaymentFailed"] end A -->|"событие"| B A -->|"событие"| B2 B -->|"событие"| C B -->|"событие"| C2 C --> A2 C2 --> B3 B3 --> A3 B2 --> A3 style A2 stroke:#5aa469,stroke-width:2px style A3 stroke:#d9534f,stroke-width:2px
Плюсы: нет единой точки отказа и единой точки изменений; сервисы связаны только контрактами событий; добавить нового участника — значит просто подписаться.
Минусы, растущие нелинейно:
- Логики саги не существует ни в одном файле. Чтобы ответить на вопрос «что происходит при отказе оплаты», нужно прочитать три сервиса и построить граф в голове.
- Циклические зависимости событий. Order слушает Payment, Payment слушает Inventory, Inventory слушает Order — и вот у вас распределённый цикл, который никто не видит.
- Нет места, где хранится состояние саги целиком. «Сколько саг сейчас висит на шаге оплаты?» — вопрос без ответа.
- Компенсации становятся отдельным набором событий, и их корректность никто не проверяет.
Хореография хороша для саг из 2–3 шагов и там, где шаги действительно независимы. На 5+ шагах она превращается в то, что Сэм Ньюмен называет «эмерджентным поведением, которое никто не проектировал».
3.2. Оркестрация: явная машина состояний
Появляется оркестратор — компонент (обычно живущий в сервисе-инициаторе), который хранит состояние саги и явно рассылает команды.
и это правильно: он подключается только после pivot
Плюсы: логика в одном месте, читается как код; состояние саги персистентно и наблюдаемо; таймауты и повторы централизованы; тестируется как обычный конечный автомат.
Минус: оркестратор — новый компонент, который легко превращается в «божественный сервис», знающий бизнес-логику чужих доменов. Лечится дисциплиной: оркестратор знает последовательность и условия перехода, но не знает, как считается НДС.
Практическое правило: до 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
Два инварианта, которые нельзя нарушать:
- Состояние саги сохраняется после каждого перехода, до отправки следующей команды. Иначе после падения оркестратор не знает, где он был.
- Компенсации не имеют права сдаться. Они повторяются бесконечно с 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"]
Три детали, которые обычно забывают и которые отличают рабочую реализацию от игрушечной:
- Хэш тела запроса. Иначе клиент, переиспользовавший ключ для другой суммы, получит ответ от старого платежа и будет уверен, что новый прошёл.
- Состояние
IN_PROGRESS. Без него два одновременных повтора выполнят операцию дважды: оба не найдут завершённой записи. - 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 сломан); - админ-интерфейс для просмотра и ручного продвижения саги: инцидент «деньги списаны, отгрузки нет» решается человеком, и ему нужен инструмент.
Типичные ошибки (по частоте, с которой встречаются в ревью):
- Компенсация не идемпотентна.
ReleaseStockвызывается дважды и возвращает товар на склад дважды. Лечится: ключ идемпотентности + inbox на стороне участника. - Нет состояния
STUCK. Сага молча уходит в DLQ после pivot, деньги списаны, никто не узнал. - Два pivot-шага. Два необратимых внешних вызова — гарантированная ручная работа при отказе второго. Проверяйте на этапе описания саги (в коде выше это делает
SagaDefinition.__init__). - Отсутствие TTL у semantic lock. Резервы накапливаются, склад «кончается» при полных полках.
- Толстый контекст саги. В
contextкладут всю корзину, документы, адреса. Таблица саг превращается во второе хранилище, миграции становятся адом. Кладите идентификаторы. - Публикация в брокер вместо outbox «потому что и так работает». Работает, пока не упадёт под нагрузкой — тогда потери обнаружатся через неделю по расхождению отчётов.
- Оркестратор знает чужую логику. Считает скидки, валидирует адреса. Возвращаемся к распределённому монолиту — см. «Монолит и модульный монолит».
- Компенсация пишет абсолютное значение вместо дельты — lost update при параллельных операциях.
9. Как выбрать механизм
Это лучший вариант — не усложняйте"] A -->|да| C["Можно ли перенести границу так,
чтобы всё оказалось в одном сервисе?"] C -->|да| D["Перенесите. Правильная граница
дешевле любой саги"] C -->|нет| E["Нужна ли строгая изоляция
промежуточных состояний?"] E -->|"да, и участники в одном кластере"| F["2PC / XA
Осознавая блокировку и цену"] E -->|нет| G["Сколько шагов?"] G -->|"2-3, компенсаций нет"| H["Хореография
на событиях"] G -->|"4+ или есть компенсации"| I["Оркестрация
явной машиной состояний"] H --> J["Обязательно:
outbox + идемпотентные получатели"] I --> J J --> K["Обязательно:
TTL у semantic lock,
состояние STUCK, алерты"] style B stroke:#5aa469,stroke-width:2px style D stroke:#5aa469,stroke-width:2px style F stroke:#c98a2b,stroke-width:2px style K stroke:#d9534f,stroke-width:2px
Обратите внимание на две зелёные ветки вверху. Лучшая распределённая транзакция — та, которой нет. Если операция постоянно требует согласованного изменения в трёх сервисах, это чаще всего сигнал о неверно проведённой границе: три сервиса на самом деле являются одним ограниченным контекстом. Прежде чем строить сагу, всерьёз рассмотрите перенос границы — материал о том, как их искать, есть в статьях «Микросервисы» и в треке по 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 доставка + идемпотентный получатель.
Источники
- Hector Garcia-Molina, Kenneth Salem. Sagas, SIGMOD 1987 — первоисточник.
- Chris Richardson. Microservices Patterns (Manning, 2018), гл. 3–4 — саги, outbox, контрмеры; онлайн-каталог: saga, transactional outbox, idempotent consumer.
- Pat Helland. Life beyond Distributed Transactions: An Apostate’s Opinion, ACM Queue.
- Jim Gray, Leslie Lamport. Consensus on Transaction Commit, 2004.
- Martin Kleppmann. Designing Data-Intensive Applications (O’Reilly, 2017), гл. 9 «Consistency and Consensus»; How to do distributed locking.
- Debezium Outbox Event Router — CDC-реализация outbox.
- Stripe: Idempotent Requests и IETF draft: The Idempotency-Key HTTP Header Field.
- Temporal documentation и AWS Step Functions Developer Guide — промышленные оркестраторы.
- Azure Architecture Center: Saga design pattern.
Что дальше
Сага и outbox описывают, как сервисы согласуют изменения между собой. Следующий вопрос — через какой интерфейс они вообще разговаривают: синхронный REST, типизированный gRPC, гибкий GraphQL или вебхуки наружу, и как менять эти контракты, не ломая потребителей.
Читайте дальше: «Стили API: REST, GraphQL, gRPC, вебхуки и версионирование».