open atlas
↑ К треку
NestJS с нуля до senior NEST · 08 · 04

Саги, компенсации и паттерн outbox

Бизнес-операция, охватывающая базы разных сервисов, не помещается в одну ACID-транзакцию, а 2PC — хрупкий SPOF. Сага сцепляет локальные транзакции компенсациями, а транзакционный outbox коммитит состояние и событие одной транзакцией БД, чтобы краш не потерял событие.

NEST Senior ◷ 19 min
Уровень
ОсновыJuniorMiddleSenior

Order-сервис сделал очевидную вещь. Внутри createOrder он вызвал await this.orders.save(order) и следующей же строкой — this.client.emit('order.created', payload). Ревью пропустило: две строки, обе зелёные в тестах, чистый деплой. Потом рутинный раскат перекатил поды, и один запрос оказался ровно между этими двумя строками в момент, когда его под получил SIGTERM. Строка заказа закоммитилась в Postgres. Событие order.created так и не дошло до брокера. Ничто не выбросило, ничто не залогировало ошибку — запрос просто испарился на полпути. Заказ существовал; инвентарь его не зарезервировал; платёж не списался. Клиент увидел «заказ оформлен», потом тишину, а через три дня тикет в поддержку спрашивал, где его покупка. В обеих строках не было бага. Баг был в зазоре между ними — две системы, никакой общей транзакции. Этот урок о том, как закрыть этот зазор, не хватаясь за распределённую транзакцию.

Почему не одна транзакция: 2PC — хрупкий SPOF

Операция «разместить заказ» затрагивает три сервиса — order, inventory, payment — и каждый владеет своей базой. Инстинкт из монолитных времён — обернуть все три записи в одну транзакцию, чтобы они коммитились или откатывались вместе. Нельзя. ACID-транзакция живёт внутри одной базы; она не может охватить три независимые базы, принадлежащие трём независимым сервисам.

Классический ответ на «атомарность поверх нескольких баз» — это двухфазный коммит (2PC / XA): координатор транзакции просит каждого участника prepare, и когда все проголосовали «да», велит всем commit. Атомарность он даёт — и почти каждый современный микросервисный стек отказывается его использовать по трём причинам, которые усугубляют друг друга под нагрузкой:

  • Координатор — единая точка отказа. Если он умирает после того, как участники проголосовали prepare, но до того, как он разослал commit, каждый участник застревает, удерживая блокировки, в ожидании решения, которое может не прийти никогда.
  • Блокировки удерживаются через сеть. Между prepare и commit каждая база держит свои блокировки на весь круг до каждого другого участника. При партишене или замедлении эти блокировки просто висят, и пропускная способность рушится.
  • Доступность падает при партишене. 2PC — это CP-протокол в жёстком смысле: если любой участник или координатор недоступен, вся транзакция блокируется. Один медленный сервис стопорит всех.

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

Сага: локальные транзакции плюс компенсации

Сага — это последовательность локальных транзакций, каждая из которых выполняется внутри собственной базы одного сервиса, связанных событиями или командами. Резервирование инвентаря — одна локальная транзакция в базе inventory. Списание платежа — другая локальная транзакция в базе payment. Каждая коммитится сама по себе. Сага — это цепочка.

Сложная часть — отказ. Если шаг 3 (списать платёж) падает, ты не можешь откатить шаги 1 и 2 — они уже закоммитились. Закоммиченная локальная транзакция исчезла; никакого ROLLBACK выпустить нельзя. Вместо этого сага запускает компенсирующие транзакции: новую локальную транзакцию, которая семантически отменяет предыдущую. Ты не «отсписываешь» платёж; ты делаешь возврат. Ты не «отрезервируешь» инвентарь; ты снимаешь резерв.

// A saga step is a local txn + a named compensation. The orchestrator runs the
// forward action; if a LATER step fails, it runs each completed step's compensate().
interface SagaStep {
  action: () => Promise<void>;      // local transaction in ONE service
  compensate: () => Promise<void>;  // semantic undo, runs AFTER commit, may itself fail
}

const placeOrderSaga: SagaStep[] = [
  { action: () => inventory.reserve(orderId, items),
    compensate: () => inventory.release(orderId) },        // not a rollback — a release
  { action: () => payment.charge(orderId, total),
    compensate: () => payment.refund(orderId) },           // not a rollback — a refund
];

Есть два способа координировать цепочку, и выбор между ними — сеньорское решение:

  • Хореография — без центрального координатора. Каждый сервис слушает события и эмитит следующее. Order-сервис эмитит order.created; inventory реагирует, резервирует, эмитит inventory.reserved; payment реагирует уже на него. Децентрализованно и просто для коротких потоков — но общий поток не живёт нигде в коде. Чтобы ответить «что происходит при размещении заказа», ты читаешь обработчики событий пяти сервисов и собираешь поток у себя в голове. Он гниёт по мере роста.
  • Оркестрация — центральный оркестратор саги владеет потоком. Он шлёт явные команды (reserveInventory, потом chargePayment), отслеживает состояние саги в своей таблице и решает, когда запускать компенсации. Это ещё один компонент, который надо построить и эксплуатировать, но поток в одном месте: явный, наблюдаемый, тестируемый, перезапускаемый.

Сеньорское эмпирическое правило: хореография для коротких потоков на 2–3 шага; оркестрация, когда поток длинный, ветвистый или требует наблюдаемости — всё, что придётся дебажить в три ночи, хочет машину состояний, которую можно запросить.

Вместе эти два стиля покрывают весь спектр координации: хореография даёт ноль новых компонентов и никакого централизованного состояния, а оркестрация — одно место, чтобы ответить «где сейчас находится эта сага?». Без этого единственного места провальный платёж в пятисервисной хореографии требует читать пять журналов событий, чтобы восстановить произошедшее.

Почему это работает

Почему оборачивание save(order) и client.emit(...) в один try/catch не делает двойную запись безопасной? Потому что два вызова бьют по двум разным системам — твоей базе и твоему брокеру сообщений — и нет границы транзакции, которая охватывала бы обе. try/catch реагирует только на исключение; он не может заставить два независимых коммита успешно-выполниться-или-упасть вместе. Краш (SIGTERM, OOM-kill, потеря питания) не выбрасывает ничего ловимого и может приземлиться между двумя вызовами в любом порядке: БД-закоммичено-но-событие-не-отправлено (событие потеряно) или событие-отправлено-но-запись-в-БД-упала (ты опубликовал ложь). Чтобы сделать состояние и событие атомарными, надо положить их в одну и ту же транзакцию базы — а событие в брокере не может быть в транзакции базы. Единственный ход — поместить событие рядом с состоянием как строку в своей же БД (outbox), закоммитить их вместе и дать отдельному relay донести эту строку до брокера потом.

Проблема двойной записи и транзакционный outbox

Инцидент из Hook — это проблема двойной записи (dual-write), и она настоящая причина существования этого урока. Сервису почти всегда приходится делать две записи как часть одного логического шага: обновить свою базу (строку заказа) и опубликовать событие, чтобы другие сервисы о нём узнали (order.created). Это две системы — БД и брокер — без общей транзакции. В каком бы порядке ты их ни делал, краш в зазоре тебя ломает:

  • Сначала БД, потом публикация, краш в промежутке → заказ существует, но событие не отправлено. Сервисы вниз по потоку никогда не узнают, что заказ существует. Инвентарь не резервирует; клиент ждёт вечно. Это Hook.
  • Сначала публикация, потом запись в БД, БД падает → ты опубликовал order.created для заказа, которого не существует. Теперь сервисы вниз по потоку действуют на ложь, и фантомный заказ расходится по системе.

Транзакционный outbox чинит это, убирая вторую систему из критической записи. В той же локальной транзакции базы, что пишет бизнес-строку, ты также делаешь INSERT события в таблицу outbox. Обе строки коммитятся вместе или ни одна — это одна транзакция в одной базе, так что зазора нет:

// Оба INSERT в ОДНОЙ локальной транзакции: бизнес-строка и строка события коммитятся атомарно.
await this.dataSource.transaction(async (manager) => {
  const order = await manager.getRepository(Order).save(newOrder);
  await manager.getRepository(OutboxEvent).insert({
    aggregateId: order.id,
    type: 'order.created',
    payload: JSON.stringify({ orderId: order.id, items: order.items }),
    // status defaults to 'pending'; a relay flips it to 'sent' after publishing
  });
}); // <- if the pod dies here, NEITHER row committed. No half state. No lost event.

Отдельный message relay затем доносит строки outbox до брокера. Две распространённые реализации:

  • Polling publisher — фоновый воркер делает SELECT строк pending, публикует каждую в брокер и помечает её sent (или удаляет). Просто, без лишней инфраструктуры; цена — задержка опроса и нагрузка.
  • Change Data Capture (CDC) — инструмент вроде Debezium читает журнал упреждающей записи (WAL) базы, видит INSERT в outbox и публикует его. Почти реальное время и без опроса, но это серьёзный кусок инфраструктуры в эксплуатации.

В любом случае relay гарантирует публикацию at-least-once: он продолжает попытки, пока брокер не подтвердит, так что отказ публикации или краш relay означает ретрай, а не потерянное событие. At-least-once означает, что консьюмер может увидеть одно и то же событие дважды — и ровно здесь зарабатывают идемпотентные консьюмеры из предыдущего урока (L03): дедупликация по id события вниз по потоку, и доставка at-least-once плюс идемпотентная обработка дают тебе effectively-once. Outbox не заменяет семантику доставки; он питает её надёжным источником.

// Relay: публикация at-least-once. Держит событие, пока брокер не подтвердит.
async function relayTick(manager: EntityManager, client: ClientProxy) {
  const pending = await manager.getRepository(OutboxEvent)
    .find({ where: { status: 'pending' }, take: 100, order: { id: 'ASC' } });
  for (const evt of pending) {
    await firstValueFrom(client.emit(evt.type, JSON.parse(evt.payload))); // may retry
    await manager.getRepository(OutboxEvent).update(evt.id, { status: 'sent' });
  }
}
Выбери лучший вариант

Размещение заказа охватывает сервисы order, inventory и payment, у каждого своя база, и должно оставаться консистентным от края до края, включая надёжную публикацию события order.created. Как ты это спроектируешь?

Викторина

Почему использовать сагу с компенсирующими транзакциями вместо одной распределённой транзакции поверх сервисов order, inventory и payment?

Викторина

Сервис делает save(order), а потом client.emit('order.created'), и под убит между двумя вызовами. Какой паттерн предотвращает потерянное событие и как?

Вспомните перед уходом
  1. 01
    Почему операция размещения заказа поверх сервисов order, inventory и payment не может использовать одну транзакцию, почему избегают 2PC и что приходит на замену?
  2. 02
    Объясни компенсации саги, хореографию против оркестрации и как транзакционный outbox решает проблему двойной записи.
Итог

Бизнес-операция, охватывающая базы, принадлежащие сервисам, не может быть одной ACID-транзакцией — транзакция живёт в одной базе — а межбазовая альтернатива 2PC/XA избегается, потому что её координатор — единая точка отказа, она держит блокировки через сеть между prepare и commit и блокируется при партишене. На замену приходит сага: последовательность локальных транзакций, каждая атомарна в одном сервисе и связана событиями или командами, принимающая eventual consistency. Поскольку закоммиченную локальную транзакцию нельзя откатить, более поздний отказ запускает компенсирующие транзакции — возврат, а не «отсписание»; снятие резерва, а не «отрезервирование» — которые выполняются после коммита, не имеют изоляции и сами могут упасть. Координируй хореографией (без координатора, сервисы реагируют на события; хорошо для 2-3 шагов, но поток не живёт нигде) для коротких потоков или оркестрацией (центральный оркестратор с явными командами и запрашиваемым состоянием), когда поток длинный, ветвистый или требует наблюдаемости. Проблема двойной записи — ядро ловушки: шаг должен закоммитить строку БД и опубликовать событие поперёк двух систем без общей транзакции, так что краш в зазоре теряет событие (Hook) или публикует ложь. Транзакционный outbox чинит это, вставляя событие в таблицу outbox в той же локальной транзакции, что и бизнес-строку — они коммитятся вместе или никак — после чего relay (polling publisher или CDC вроде Debezium) доносит строки outbox до брокера at-least-once. В сочетании с идемпотентными консьюмерами из L03, дедуплицирующими по id события, публикация at-least-once плюс идемпотентная обработка дают effectively-once доставку, при которой ни одна запись не сбегает из транзакции и ни одно событие не теряется из-за краша. Теперь, когда встретишь тикет «клиент видит “заказ оформлен”, но ничего не происходит», будешь знать, куда смотреть: зазор двойной записи без outbox или шаг саги, так и не запустивший компенсацию.

Практика

Начни сверху. Задачи идут от простого к сложному: вспомнить факт, применить к случаю, затем senior-уровень. Открой, попробуй, потом открой ответ.

вспомнитьприменитьуглубить0 из 5 завершено

Что-то непонятно?

Задай вопрос по этому уроку. Вопросы анонимны и попадают напрямую автору — урок станет лучше.

хоткеи развернуть
поиск
K
пред. пьеса
k
след. пьеса
j
тиры
t
это меню
?
sources2
expand
  1. 01
  2. 02

Trademarks belong to their respective owners. Editorial reference only.