Распределённые транзакции: 2PC, saga, outbox, компенсации
Локальная транзакция — это соглашение с одним ресурсом: «либо всё, либо ничего, и я тебе честно скажу, что получилось». Соглашение держится на одном физическом факте: у базы есть один журнал, и в нём есть одна запись коммита, которая либо сброшена на диск, либо нет. Атомарность — это не магия менеджера транзакций, это атомарность одного fsync.
Как только эффект операции размазан по двум ресурсам — заказ в PostgreSQL и списание в биллинге, строка в базе и сообщение в Kafka, — единственной записи коммита больше нет. Есть две, и между ними существует момент, в который процесс может умереть. Вся эта статья — про то, что делают с этим моментом.
Три ответа, между которыми выбирают на практике:
- Договориться о едином решении — 2PC и его потомки. Атомарность настоящая, цена — блокировки и уязвимость к отказу координатора.
- Отказаться от атомарности и научиться откатывать эффекты бизнес-действиями — сага и компенсации. Цена — отсутствие изоляции и видимые промежуточные состояния.
- Свести две записи к одной — outbox. Цена — отложенная публикация и гарантированные дубликаты.
Дальше предполагается, что вы уже прочли про модели отказов (главное: «упал» и «медленно отвечает» неразличимы), про модели согласованности и хотя бы бегло — про консенсус. Архитектурный взгляд на сагу как на стиль интеграции сервисов разобран отдельно, в «Saga и распределённые транзакции»; здесь нас интересует фундамент: какие протоколы что доказуемо гарантируют, где именно они ломаются и как этот излом выглядит в логах.
Задача атомарного коммита — это не консенсус
Формально задача называется atomic commitment (AC). Есть N участников, каждый голосует yes или no. Требуется:
- AC1 (соглашение). Все участники, принявшие решение, приняли одно и то же.
- AC2 (валидность).
commitвозможен, только если все проголосовалиyes. - AC3 (нетривиальность). Если все проголосовали
yesи не было отказов, решение —commit. - AC4 (завершаемость). Каждый исправный участник в конце концов принимает решение.
Похоже на консенсус, но условие AC2 меняет всё. В консенсусе достаточно, чтобы решение было чьим-то предложением, — большинство может решить за всех. В атомарном коммите один-единственный no обязан победить любое большинство: если склад не смог зарезервировать товар, никакие четыре других сервиса не имеют права закоммитить. Отсюда прямое следствие: атомарный коммит не переживает отказ произвольного участника так, как консенсус переживает отказ меньшинства. Голос участника нельзя заменить голосом соседа.
Именно поэтому теорема FLP (Fischer, Lynch, Paterson, 1985) бьёт здесь особенно больно: в асинхронной модели с одним возможным отказом нельзя одновременно гарантировать AC1 и AC4. Реальные протоколы жертвуют AC4 — то есть зависают.
Двойная запись: самая частая распределённая транзакция, которую никто так не называет
Прежде чем разбирать 2PC, стоит увидеть, что распределённая транзакция уже есть в вашем коде, даже если вы никогда не писали XA START.
def create_order(cmd):
with db.transaction(): # граница атомарности №1
order = orders.insert(cmd)
broker.publish("orders.created", order) # граница атомарности №2
return order
Между commit и publish процесс может быть убит OOM-killer’ом, кубернетесом при rolling update или уйти в 40-секундную STW-паузу GC, после которой брокер уже разорвал соединение. Каждый из четырёх исходов на схеме — отдельный жанр инцидента, и три из них обнаруживаются не в логах, а сверкой множеств: «заказ есть, события нет» ищется как разность count(orders) и числа сообщений в топике, а не как строка ERROR.
Сценарий отказа целиком. Сервис orders пишет заказ и публикует событие. Деплой в 14:03, поды перезапускаются с terminationGracePeriodSeconds: 30, а linger.ms продюсера — 100 мс, но в буфере 8 МБ неотправленных батчей. Под получает SIGKILL. В логе приложения — ничего: последняя строка order 4711 created, потому что publish асинхронный и колбэк не успел выполниться. В БД заказ есть. В топике его нет. Через час склад присылает суточную сверку, где 217 заказов «не существуют». Дедукция занимает день, потому что отсутствие сообщения не оставляет следов — искать надо разницу двух множеств, а не строку в логе.
Ключевой вывод, к которому мы вернёмся в разделе про outbox: проблема не в брокере и не в базе. Проблема в том, что в коде две границы атомарности, а сценарий требует одной.
Two-Phase Commit
Протокол
2PC описан Джимом Греем ещё в «Notes on Data Base Operating Systems» (1978) и с тех пор не изменился по сути.
Фаза 1 (голосование). Координатор рассылает PREPARE. Каждый участник делает всё, кроме коммита: проверяет ограничения, берёт блокировки, принудительно сбрасывает на диск запись prepared — и отвечает yes или no. Отвечая yes, участник даёт обещание: «я смогу закоммитить в любой момент в будущем, даже если сейчас упаду и перезагружусь». Это обещание — центральный элемент протокола и источник всех бед.
Фаза 2 (решение). Если все ответили yes, координатор принудительно записывает в свой журнал COMMIT — вот точка коммита, момент, после которого транзакция считается совершившейся, — и рассылает COMMIT. Участники применяют, отпускают блокировки, отвечают ACK.
горизонт xmin не двигается,
ждём координатора неограниченно A-->>TC: могу ли я узнать решение? B-->>TC: могу ли я узнать решение?
Стоимость на успешном пути: 4(N−1) сообщений (prepare, vote, decision, ack), 2 сетевых RTT до ответа клиенту, N+1 принудительных сбросов журнала — по одному у каждого участника плюс решение координатора. Оптимизация presumed abort (Mohan, Lindsay, Obermarck, TODS 1986) убирает часть записей: отсутствие информации о транзакции трактуется как abort, поэтому запись об аборте не нужна вовсе, а ACK для аборта можно не ждать.
Где именно он блокируется
Между отправкой yes и получением решения участник находится в состоянии неопределённости (uncertainty period). Он уже не может откатиться сам (обещал), но ещё не знает, коммитить ли. Всё это время блокировки удерживаются.
(единственный законный переход) Prepared --> Heuristic: оператор решает вручную Heuristic --> [*]: возможно нарушение атомарности Committed --> [*] Aborted --> [*] note right of Prepared Ресурсы удержаны. Восстановление после падения возвращает узел СЮДА ЖЕ, а не в Aborted. end note
Обратите внимание на петлю Prepared --> Prepared. Это не украшение диаграммы, это теорема. Skeen и Stonebraker (1983) доказали: в асинхронной модели с возможностью отказа коммуникации не существует неблокирующего протокола атомарного коммита. Участник в prepared не может решить сам, потому что:
- закоммитить нельзя — вдруг кто-то другой проголосовал
no, и решение былоabort; - откатиться нельзя — вдруг все проголосовали
yes, координатор записалcommit, а сосед уже применил его и отдал результат клиенту.
Спросить соседей тоже не всегда помогает: если сосед сам в prepared, он знает не больше. Это называется кооперативным терминационным протоколом, и он разрешает ситуацию только тогда, когда хоть кто-то уже получил решение.
Как это выглядит в PostgreSQL
PostgreSQL умеет быть участником 2PC: PREPARE TRANSACTION 'gid', затем COMMIT PREPARED 'gid' или ROLLBACK PREPARED 'gid'. По умолчанию это выключено — max_prepared_transactions = 0, и это осознанное решение разработчиков.
Сценарий отказа. Координатор (сервер приложений на Narayana) упал между PREPARE TRANSACTION и COMMIT PREPARED, его переналили в другую зону — с пустым журналом транзакций, потому что журнал лежал в emptyDir. Что происходит дальше:
-- в pg_stat_activity пусто: сессии нет, приложение не подключено
SELECT gid, prepared, owner, database FROM pg_prepared_xacts;
-- gid | prepared | owner | database
-- -------------------+-------------------------------+-------+----------
-- 131077_orders_4711| 2026-05-11 02:14:07.113322+03 | app | shop
-- блокировки живут, но принадлежат «виртуальной транзакции» -1/NNN
SELECT locktype, relation::regclass, mode, virtualtransaction
FROM pg_locks WHERE virtualtransaction LIKE '-1/%';
А в логе сервера каждые несколько минут одно и то же:
LOG: automatic vacuum of table "shop.public.orders": index scans: 0
tuples: 0 removed, 12844411 remain, 8921033 are dead but not yet removable,
oldest xmin: 918274612
dead but not yet removable при неподвижном oldest xmin — это подпись зависшей prepared-транзакции (или зависшего репликационного слота, или очень долгого REPEATABLE READ). Диск растёт, планы деградируют, и через несколько дней вы упираетесь в wraparound-защиту:
WARNING: database "shop" must be vacuumed within 10985430 transactions
HINT: To avoid a database shutdown, execute a database-wide VACUUM in that database.
Так одна забытая строка в pg_prepared_xacts останавливает продакшн. Практический вывод: если вы включаете max_prepared_transactions, вы обязаны завести алерт pg_prepared_xacts с prepared < now() - interval '5 minutes' — раньше, чем алерт на место на диске.
XA, эвристические решения и честный запасной выход
Промышленный стандарт интерфейса участника — X/Open XA. У него есть операция, которой нет в учебниках по 2PC: эвристическое решение. Когда участник висит в prepared слишком долго, менеджер ресурсов (или оператор) может решить самостоятельно, приняв риск нарушить атомарность. Исход обязан быть зафиксирован и сообщён:
| Код | Что означает |
|---|---|
XA_HEURCOM |
участник самовольно закоммитил |
XA_HEURRB |
участник самовольно откатился |
XA_HEURMIX |
часть работы закоммичена, часть откачена — атомарность нарушена точно |
XA_HEURHAZ |
исход неизвестен даже участнику |
В логе менеджера транзакций (Narayana) это выглядит примерно так:
WARN [com.arjuna.ats.arjuna] ARJUNA012117: TransactionReaper::check timeout for TX
0:ffffc0a80105:5f3a:6819b2c1:2f in state RUN
WARN [com.arjuna.ats.jta] Received heuristic outcome XAException.XA_HEURHAZ
from resource ledger-db for xid 0:ffffc0a80105:5f3a:6819b2c1:2f
ERROR [billing] Transaction outcome UNKNOWN for order 4711; manual reconciliation required
После этого требуется xa_forget и — обязательно — человек со сверкой. XA_HEURMIX в логе означает, что деньги списаны, а товар не зарезервирован (или наоборот), и никакой автоматики, которая это исправит, в протоколе нет. В MySQL те же зависшие ветви ищутся командой XA RECOVER, в PostgreSQL — через pg_prepared_xacts; и там, и там разгребает человек.
Эвристики — это признание проектировщиков стандарта: неблокирующего 2PC не бывает, поэтому вот вам рубильник и ответственность за него.
3PC и почему он не спасает
Skeen (1981) предложил трёхфазный коммит: между голосованием и решением добавляется фаза pre-commit, которая сообщает участникам «все проголосовали yes». Тогда участник, оказавшийся без координатора, но уже получивший pre-commit, может смело коммитить, а не получивший — смело откатывать.
Ловушка в модели отказов. 3PC неблокирующий только при синхронной сети с известной верхней границей задержки и отказах типа fail-stop. При разделении сети он ломается прямо и грубо: одна половина кластера получила pre-commit и коммитит по таймауту, вторая не получила и откатывает по таймауту. Получаем не блокировку, а нарушение атомарности — то есть строго худший исход. Плюс третья фаза добавляет ещё один RTT к и без того дорогому протоколу. Именно поэтому 3PC не используется практически нигде; вместо него берут реплицированный координатор.
2PC — это Paxos Commit с одним акцептором
Самое полезное, что можно понять про 2PC, сформулировали Джим Грей и Лесли Лампорт в «Consensus on Transaction Commit» (TODS, 2006). Их результат: 2PC — это вырожденный частный случай Paxos Commit, в котором каждое голосование решается одним-единственным акцептором — координатором.
Отсюда лекарство очевидно: сделайте решение о коммите реплицированным консенсусом вместо записи в журнал одного узла. Тогда падение координатора перестаёт быть фатальным — новый лидер узнает решение из реплицированного лога. Paxos Commit с 2F+1 акцепторами переживает F отказов и стоит столько же RTT, сколько 2PC (при совмещении фаз), но требует больше сообщений.
Важная оговорка: это лечит отказ координатора, но не отказ участника. Участник, который умер в prepared навсегда, блокирует транзакцию в любом протоколе — потому что только он держит свои данные. Лечится это репликацией самого участника, что и делают все системы ниже.
Где 2PC жив и прекрасно работает
Тезис «2PC умер» — неверен. Умер XA поверх нереплицированных ресурсов и разнородных вендоров. Внутри современных распределённых СУБД 2PC — рабочая лошадь, потому что там оба лекарства применены сразу.
| Система | Как устроено | Что чинит проблему координатора |
|---|---|---|
| Spanner | 2PC поверх Paxos-групп; каждый участник — Paxos-группа, координатор — лидер одной из них | состояние координатора реплицировано; новый лидер продолжает с того же места; TrueTime даёт внешнюю согласованность (OSDI 2012) |
| Percolator / TiDB | 2PC поверх BigTable/TiKV; «точка коммита» — запись в первичную строку блокировки | решение хранится в реплицированном KV, а не в памяти координатора; клиент может умереть, любой другой клиент дочистит (OSDI 2010) |
| CockroachDB | запись транзакции (transaction record) в Raft-группе + write intents; parallel commits | запись транзакции реплицирована; чужие транзакции, наткнувшись на intent, «толкают» её и дочищают |
| MongoDB (шардированные транзакции) | 2PC с координатором на одном из шардов | координатор пишет решение в реплицированную коллекцию config.transaction_coordinators |
| Kafka (транзакции продюсера) | transaction coordinator, состояние в реплицированном топике __transaction_state, маркеры COMMIT/ABORT в партициях данных |
координатор переизбирается вместе с лидером партиции __transaction_state |
Общий рецепт: координатор не должен быть единственной копией решения. Если ваш «менеджер распределённых транзакций» — это Java-процесс с журналом в локальной файловой системе пода, у вас 2PC 1978 года со всеми его свойствами.
Saga: обмен отката на компенсацию
Формальная модель
Сага придумана не для микросервисов. Garcia-Molina и Salem (SIGMOD 1987) решали проблему долгоживущих транзакций внутри одной СУБД: транзакция, которая держит блокировки часами, убивает конкурентность, поэтому её разбивают на короткие транзакции T₁…Tₙ, каждая из которых коммитится сразу. Гарантия саги: система в конце концов оказывается либо в состоянии «выполнены все T₁…Tₙ», либо «выполнены T₁…T_j, а затем компенсации C_j…C₁». Никаких обещаний про то, что происходит в середине, сага не даёт. Формально это ACD: атомарность (в смысле «всё-или-семантически-ничего»), согласованность, долговечность — но без изоляции.
Компенсация — это не откат
| Откат (rollback) | Компенсация | |
|---|---|---|
| Кто видел эффект | никто | все |
| Восстанавливает | точное прежнее состояние | семантически приемлемое состояние |
| Пример | UNDO записи в WAL |
«возврат средств», а не «отмена списания» |
| Следы в истории | не остаётся | остаётся навсегда: две строки в выписке, письмо клиенту, запись в аудите |
| Может не сработать | нет | да, и это нормальный сценарий |
Компенсация должна быть идемпотентной (её будут ретраить) и коммутативно-безопасной (она может прийти после того, как клиент уже увидел эффект). А главное — компенсация не всегда существует: письмо отправлено, SMS доставлена, товар выдан курьеру, платёж ушёл контрагенту по SWIFT. Отсюда правило порядка шагов, которое и делает сагу реализуемой: сначала компенсируемые (их можно отменить), затем ровно один pivot — точка невозврата, последний откатываемый или первый неоткатываемый шаг, — затем повторяемые, которые гарантированно завершатся при достаточном числе ретраев («отправить письмо» можно ретраить вечно).
Жизненный цикл и оркестратор
компенсируемый"] R1 -->|ок| R2["Шаг 2: захолдировать деньги
компенсируемый"] R1 -->|ошибка| F(["Провалено: компенсаций не требуется"]) R2 -->|ок| P["Шаг 3: списать деньги
PIVOT — точки возврата больше нет"] R2 -->|ошибка| C1["C1: снять резерв"] P -->|ок| R4["Шаг 4: создать задание на отгрузку
повторяемый"] P -->|ошибка| C2["C2: снять холд"] C2 --> C1 C1 --> F R4 -->|ошибка| RT["Ретрай с backoff
вечно, до успеха"] RT --> R4 R4 -->|ок| D(["Завершено"]) style P stroke:#d88a6f,stroke-width:2px style RT stroke:#b48ead,stroke-width:2px
Псевдокод оркестратора с восстановлением после падения:
выполнить(saga):
(позиция, направление) ← журнал.прочитать(saga.id) # источник истины, не память процесса
если направление = НАЗАД: откатить(saga, с=позиция); вернуть ПРОВАЛЕНО
для i от позиция до n:
журнал.записать(saga.id, шаг=i, НАЧАТ) # ДО вызова, иначе потеряем факт вызова
результат ← вызвать(шаги[i], ключ=(saga.id, i))
если результат ≠ УСПЕХ и i ≤ pivot:
журнал.записать(saga.id, направление=НАЗАД, позиция=i-1)
откатить(saga, с=i-1); вернуть ПРОВАЛЕНО
если результат ≠ УСПЕХ:
повторять шаг i с экспоненциальным backoff # после pivot назад дороги нет
журнал.записать(saga.id, шаг=i, ГОТОВ)
вернуть УСПЕХ
Ключевая строка — запись НАЧАТ до сетевого вызова. Если писать после, то падение между вызовом и записью приведёт к тому, что при восстановлении шаг будет выполнен повторно, а оркестратор будет считать, что не выполнял его вовсе. Отсюда следует, что каждый шаг обязан быть идемпотентным по ключу (saga_id, step) — этому посвящена следующая статья.
Реализация на Python — минимальный, но честный оркестратор:
import time
from dataclasses import dataclass
from typing import Callable, Sequence
Action = Callable[[dict, str], None] # (контекст, ключ идемпотентности)
@dataclass(frozen=True)
class Step:
name: str
action: Action
compensation: Action | None # None — компенсировать нечего (чтение, логирование)
pivot: bool = False # после pivot откат невозможен
def run_saga(saga_id: str, steps: Sequence[Step], ctx: dict, log) -> bool:
"""log — долговечный журнал в той же БД, что и состояние оркестратора."""
start, direction = log.load(saga_id) # возобновление после падения процесса
if direction == "BACKWARD":
return _compensate(saga_id, steps, ctx, log, frm=start)
pivot = next((i for i, s in enumerate(steps) if s.pivot), len(steps))
for i in range(start, len(steps)):
log.append(saga_id, i, "STARTED") # сначала намерение, потом действие
try:
steps[i].action(ctx, f"{saga_id}:{i}") # ключ стабилен между всеми ретраями
except Exception:
if i <= pivot: # ещё можно назад
log.append(saga_id, i, "FAILED")
return _compensate(saga_id, steps, ctx, log, frm=i - 1)
_retry_forever(steps[i].action, ctx, f"{saga_id}:{i}")
log.append(saga_id, i, "DONE")
return True
def _compensate(saga_id, steps, ctx, log, frm: int) -> bool:
for i in range(frm, -1, -1): # строго в обратном порядке
if steps[i].compensation is None:
continue
log.append(saga_id, i, "COMPENSATING")
_retry_forever(steps[i].compensation, ctx, f"{saga_id}:{i}:comp")
log.append(saga_id, i, "COMPENSATED")
return False
def _retry_forever(fn: Action, ctx: dict, key: str, cap: float = 60.0) -> None:
"""Компенсация не имеет права сдаться: выход из цикла — успех или человек."""
delay, attempt = 0.5, 0
while True:
try:
return fn(ctx, key)
except Exception:
attempt += 1
if attempt == 10: # алерт с бизнес-контекстом, не стектрейс
escalate(f"шаг {key} не проходит 10 раз; нужна ручная сверка")
time.sleep(min(delay, cap))
delay *= 2
Сложность. Успешный путь: n сетевых вызовов и 2n записей в журнал — O(n) по времени, O(n) по памяти на журнал (журнал обязан быть долговечным и хранить не меньше, чем максимальный горизонт ретраев). Худший случай — провал на шаге n−1: n−1 вызовов вперёд плюс n−2 компенсаций, то есть не более 2n вызовов, O(n). Латентность саги — сумма латентностей шагов, а не максимум: сага принципиально последовательна там, где 2PC параллелен. Если шаги независимы, их можно распараллелить, но тогда компенсации тоже параллельны и порядок компенсации перестаёт быть строго обратным — это допустимо только если компенсации коммутируют.
Изоляции нет: аномалии и контрмеры
Самая дорогая ошибка в сагах — забыть, что буква «I» из ACID выпала. Промежуточные состояния видны всем.
| Аномалия | Что происходит | Пример |
|---|---|---|
| Грязное чтение | другая транзакция видит эффект шага, который потом будет компенсирован | отчёт по выручке посчитал платёж, который через 3 секунды вернули |
| Потерянное обновление | другая сага меняет ту же сущность между шагами | два заказа резервируют последний товар, оба «успешны» |
| Нечёткое чтение | сага дважды читает данные и видит разное | проверка лимита на шаге 1 прошла, к шагу 3 лимит исчерпан |
Контрмеры (систематизированы Крисом Ричардсоном в «Microservices Patterns», восходят к работам по семантическим ACID-свойствам):
- Семантическая блокировка. Сущность получает флаг
state = PENDING, и другие саги обязаны его уважать. По сути вы вручную реализуете блокировку, которую вам не дал менеджер транзакций, — с тем же риском взаимоблокировок и с обязательным TTL. - Коммутативные обновления.
balance = balance - 100вместоbalance = 900; компенсация+100тогда коммутирует с параллельными операциями. Это ровно та же идея, что в CRDT-счётчиках из статьи про репликацию. - Пессимистичное представление. Меняем порядок шагов так, чтобы риск грязного чтения был минимальным: сначала списываем, потом начисляем, а не наоборот.
- Перечитывание значения. Перед записью перечитать и сверить версию — оптимистическая блокировка; при расхождении сага провалилась.
- Файл версий. Записываем операции как события и переупорядочиваем их при применении — превращаем некоммутативные операции в коммутативные.
- По значению. Выбирать механизм под транзакцию: дешёвые и низкорисковые — сагой, критичные (перевод между счетами внутри одного банка) — обычной ACID-транзакцией. Самая недооценённая контрмера: не делать распределённую транзакцию там, где можно не делать.
Сценарий отказа. Сага бронирования: шаг 1 — «зарезервировать место», шаг 2 — «списать оплату». Между ними 800 мс. Резерв не помечен PENDING, потому что «так проще». Два клиента параллельно бронируют последнее место: обе саги проходят шаг 1 (проверка free > 0 и декремент в разных транзакциях), обе списывают деньги. В логах — два безупречных saga completed. Проблема всплывает на стойке регистрации. Ни одна метрика не покажет этого, потому что ошибки не было: сага сработала ровно так, как написана.
Что делать, когда упала компенсация
Компенсация может провалиться навсегда: сервис удалён, деньги ушли контрагенту, товар отгружен. Тогда:
- Ретраить с backoff бесконечно — компенсация не имеет права «сдаться» в коде.
- После N попыток — алерт с бизнес-контекстом, а не
NullPointerException: «саге order-4711 требуется ручной возврат 12 400 RUB». - Персистентная очередь ручных задач (
saga_manual_intervention), а не DLQ, который никто не читает. - Метрика
saga_stuck_totalпо типам саг — это операционный SLI, и он должен быть на дашборде рядом с latency.
Отсутствие пункта 3 — самая частая архитектурная дыра в сагах: код предусматривает компенсацию, но не предусматривает провал компенсации, поэтому падает в catch (Exception e) { log.error(...) }, и деньги зависают навсегда.
Transactional outbox: как свести две записи к одной
Идея
Вернёмся к двойной записи. Атомарно накрыть БД и брокер может только 2PC — а брокеры XA либо не поддерживают, либо поддерживают так, что лучше бы не поддерживали. Outbox решает задачу иначе: делает публикацию частью той же самой транзакции БД, записывая событие в обычную таблицу.
BEGIN;
INSERT INTO orders (id, customer_id, total, status)
VALUES ('4711', 'c-42', 12400, 'CREATED');
INSERT INTO outbox (id, aggregate_type, aggregate_id, event_type, payload)
VALUES (gen_random_uuid(), 'Order', '4711', 'OrderCreated',
jsonb_build_object('orderId', '4711', 'total', 12400));
COMMIT; -- одна граница атомарности: либо и заказ, и событие, либо ничего
Отдельный процесс (relay) читает outbox и публикует в брокер. Состояние «заказ есть, события нет» стало невозможным — событие лежит в той же базе и будет опубликовано. Состояние «событие есть, заказа нет» тоже невозможно. Взамен появляется другое: событие может быть опубликовано дважды, потому что relay может упасть между publish и отметкой об отправке.
Relay: polling против CDC
Polling. Простой воркер читает неотправленные строки и публикует.
-- порядок по created_at/id и блокировка «своих» строк
WITH batch AS (
SELECT id FROM outbox
WHERE published_at IS NULL
ORDER BY created_at
LIMIT 100
FOR UPDATE SKIP LOCKED
)
UPDATE outbox SET published_at = now()
FROM batch WHERE outbox.id = batch.id
RETURNING outbox.*;
Плюсы: никакой инфраструктуры, работает в любой БД. Минусы: задержка = интервал опроса, нагрузка на БД от постоянного сканирования, необходимость чистить таблицу (партиционирование по дням + DROP PARTITION, иначе outbox за полгода станет крупнейшей таблицей в кластере).
CDC. Debezium читает WAL/binlog и публикует события с трансформацией io.debezium.transforms.outbox.EventRouter, которая превращает строку outbox в сообщение с нужным топиком и ключом. Плюсы: нулевая нагрузка на БД сверх репликации, задержка в единицы миллисекунд, порядок WAL = порядок коммитов (это важно, см. ниже). Минусы: Kafka Connect в эксплуатации, слот репликации как новая точка отказа.
Три способа сломать outbox
1. Опрос по монотонному id теряет строки. bigserial выдаёт номер при INSERT, а не при COMMIT. Транзакция A получила id = 100, транзакция B получила id = 101 и закоммитилась первой. Relay читает WHERE id > 99, видит только 101, запоминает last_seen = 101. Через 5 мс коммитится A со id = 100 — и её событие не будет прочитано никогда. Ошибки нет, лога нет, метрика «размер очереди» нулевая; обнаруживается сверкой через недели. Лечится тремя способами: не использовать монотонный курсор (фильтровать по published_at IS NULL, как в запросе выше); использовать CDC, где порядок задан журналом коммитов; либо читать только строки ниже xmin снапшота pg_current_snapshot().
2. Несколько экземпляров relay ломают порядок. FOR UPDATE SKIP LOCKED с четырьмя воркерами отлично масштабируется и одновременно гарантирует, что OrderUpdated может уйти в брокер раньше OrderCreated. Если порядок важен (а он важен почти всегда), шардируйте по aggregate_id: воркер обрабатывает только те строки, у которых hashtext(aggregate_id) % N = worker_id. Порядок сохраняется в пределах агрегата — ровно та гарантия, которую даёт и партиция Kafka.
3. Репликационный слот сжирает диск. Debezium остановлен на выходные (упал коннектор, никто не заметил). PostgreSQL обязан хранить WAL, пока слот его не прочитал:
SELECT slot_name, active, wal_status,
pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)) AS retained
FROM pg_replication_slots;
-- slot_name | active | wal_status | retained
-- ---------------+--------+------------+----------
-- debezium_shop | f | extended | 412 GB
Финал предсказуем:
PANIC: could not write to file "pg_wal/xlogtemp.31415": No space left on device
LOG: server process (PID 31415) was terminated by signal 6: Aborted
LOG: database system is shut down
Неактивный слот роняет основную базу, а не CDC-пайплайн. Обязательные меры: max_slot_wal_keep_size (с PG 13 слот перейдёт в wal_status = 'lost' и умрёт вместо базы), алерт на active = false и на retained больше порога.
Inbox: вторая половина решения
Outbox даёт at-least-once. Значит, у потребителя обязан быть дедуп — inbox. И критически важно, где именно он коммитится:
BEGIN;
-- дедуп и эффект — в одной транзакции; иначе вы снова вернулись к двойной записи
INSERT INTO inbox (consumer_group, message_id, processed_at)
VALUES ('billing', '2b1a...-c4', now())
ON CONFLICT DO NOTHING;
-- 0 строк вставлено → сообщение уже обработано → выходим и коммитим ACK
UPDATE accounts SET balance = balance - 12400 WHERE id = 'c-42';
COMMIT;
Если отметка о дедупе и эффект коммитятся раздельно — вы построили распределённую транзакцию с двумя участниками и вернулись в начало статьи.
Варианты и альтернативы
- Transactional messaging в брокере. RocketMQ умеет «half messages»: продюсер отправляет невидимое сообщение, коммитит локальную транзакцию, затем подтверждает; при неопределённости брокер сам опрашивает продюсера. По сути тот же 2PC, где вторым участником выступает приложение. Родственная мелочь:
NOTIFY/LISTENв PostgreSQL уменьшает задержку relay, но уведомления не долговечны — периодический опрос всё равно обязателен. - Event sourcing. Если состояние и есть лог событий, outbox не нужен: публиковать нечего, кроме того, что уже записано. См. CQRS и event sourcing.
- TCC (Try-Confirm/Cancel). Явные резервирования вместо блокировок:
tryсоздаёт резерв с TTL,confirmего материализует,cancelотпускает. Это сага, у которой первая фаза похожа наprepare, но блокировки заменены на бизнес-сущность «резерв» с собственным сроком годности. Хорошо ложится на домены, где резерв — естественное понятие (билеты, склад, лимиты).
Почему exactly-once обычно миф
Формулировка, которую стоит запомнить дословно: exactly-once delivery невозможна; exactly-once effect достижим.
Невозможность доставки — это проблема двух генералов (Akkoyunlu, Ekanadham, Huber, 1975; название дал Джим Грей). Отправитель послал сообщение и не получил подтверждения; отличить «сообщение потерялось» от «сообщение дошло, потерялся ACK» он не может. Выбор ровно из двух вариантов: не ретраить — at-most-once, сообщение может пропасть; ретраить — at-least-once, сообщение может продублироваться. Третьего не дано никаким протоколом, потому что никакое конечное число раундов не устраняет неопределённость последнего сообщения. Это не инженерное несовершенство, а доказанное свойство канала с потерями.
Что тогда означает «exactly-once semantics» в маркетинге Kafka и Flink? Строго следующее: дедупликация и фиксация прогресса выполняются атомарно с самим эффектом, внутри системы, которая контролирует и то, и другое.
- Kafka EOS. Идемпотентный продюсер даёт
(PID, epoch, sequence)на каждую партицию, и брокер отбрасывает повторы внутри окна. Транзакции продюсера атомарно коммитят «сообщения в топики + смещения консьюмера» через маркеры COMMIT/ABORT и топик__transaction_state. Это работает строго для контура Kafka → обработка → Kafka. Как только у обработки есть побочный эффект вовне — HTTP-вызов,INSERTв чужую БД, отправка письма — гарантия заканчивается на границе Kafka, и exactly-once обязан обеспечивать приёмник. - Flink. Барьерные снапшоты Чанди–Лампорта плюс
TwoPhaseCommitSinkFunction: сток предварительно записывает данные (preCommit), состояние чекпоинта фиксируется, затем сток коммитит. Это снова 2PC, где вторым участником назначен внешний приёмник — и он обязан уметьprepare/commit(транзакционная БД умеет, «отправить SMS» — нет).
Сценарий отказа, показывающий цену «exactly-once» в Kafka. Приложение с transactional.id упало, не закрыв транзакцию; transaction.timeout.ms выставлен в 15 минут «чтобы не падало на длинных батчах». Потребители с isolation.level=read_committed не могут читать дальше last stable offset:
# лаг растёт линейно, при этом данные в топик пишутся
kafka.log:type=Log,name=LogEndOffset,topic=payments,partition=7 → 8 941 220
kafka.log:type=Log,name=LastStableOffset,topic=payments,partition=7 → 8 827 311 (стоит 40 минут)
$ kafka-transactions.sh --bootstrap-server kafka-1:9092 find-hanging --topic payments
Topic Partition ProducerId ProducerEpoch StartOffset LastTimestamp Duration(min)
payments 7 21004 3 8827311 2026-05-11T02:14:07Z 214
Это тот же самый паттерн, что и pg_prepared_xacts: незакрытая транзакция блокирует прогресс всех остальных. Разница лишь в том, что здесь есть встроенный таймаут — а в XA его нет. Лечение: kafka-transactions.sh abort --topic payments --partition 7 --start-offset 8827311, а профилактика — разумный transaction.timeout.ms и алерт на разрыв LEO и LSO.
Что делают вместо exactly-once. Всегда одно и то же:
- Транспорт — at-least-once (ретраи включены,
acks=all). - У каждого сообщения — стабильный идентификатор, порождённый источником, а не транспортом:
message_idиз outbox, а не «UUID, сгенерированный при отправке» (иначе ретрай создаст новый id, и дедуп не сработает). - Обработчик идемпотентен, а отметка о дедупе коммитится в той же транзакции, что и эффект.
- Окно дедупликации больше максимального горизонта ретраев с запасом — иначе поздний дубль пройдёт мимо.
Подробный разбор — в следующей статье трека.
Как выбирать
одна и та же БД?"} B -->|"да"| C["Обычная ACID-транзакция.
Не изобретайте распределённость"] B -->|"нет"| D{"Второе место —
брокер сообщений?"} D -->|"да"| E["Transactional outbox
+ inbox у потребителя"] D -->|"нет"| F{"Промежуточное состояние
допустимо видеть снаружи?"} F -->|"нет, нужна изоляция"| G{"Ресурсы реплицированы
и под общим контролем?"} G -->|"да: Spanner, TiDB,
CockroachDB, YDB"| H["Встроенный 2PC поверх консенсуса"] G -->|"нет: разные вендоры,
нереплицированный координатор"| I["XA. Считайте цену:
блокировки, эвристики,
ручная сверка"] F -->|"да"| J{"Каждый шаг
компенсируем?"} J -->|"да"| K["Сага: оркестрация
+ семантические блокировки"] J -->|"нет"| L["Переставить шаги:
некомпенсируемые после pivot,
или заменить на резервирование TCC"] style C stroke:#6f9fd8,stroke-width:2px style I stroke:#d88a6f,stroke-width:2px
Три оси, по которым различаются варианты: изоляция (есть у 2PC, нет у саги), латентность (2 RTT у 2PC против суммы шагов у саги — сага принципиально последовательна) и операционный риск (зависшие prepared-транзакции у XA, застрявшие компенсации у саги, растущие таблица и слот у outbox). Четвёртая, неформальная ось — состав участников: XA требует ресурсов с поддержкой XA, встроенный 2PC — узлов одной СУБД, а сага работает с чем угодно, включая внешний HTTP-API платёжного шлюза, у которого никакого prepare нет и не будет.
Типичные ошибки
- Считать, что двойная запись «почти всегда работает». При 10 000 заказов в сутки и вероятности сбоя 0,1 % это 10 расхождений в день — то есть 3650 в год, каждое из которых кто-то разбирает руками.
- Включить
max_prepared_transactionsбез алерта на зависшие GID. Через неделю вы будете объяснять, почемуVACUUMне работает и база уходит в защиту от wraparound. - Держать журнал координатора в
emptyDirпода. Журнал координатора — это единственная копия решения о коммите. Потеря его равна потере атомарности. - Считать компенсацию откатом. Она видима, она может не сработать, и она обязана быть идемпотентной. И её обязательно надо тестировать — компенсации ломаются чаще основных шагов, потому что выполняются в тысячу раз реже.
- Забыть про изоляцию в саге. «Мы же сделали сагу правильно» не спасает от двух броней на последнее место. Семантические блокировки или коммутативные операции — обязательны, а не опциональны.
- Не проектировать провал компенсации. Нужны бесконечные ретраи, алерт с бизнес-контекстом и очередь ручных задач, а не
log.error. - Опрашивать outbox по монотонному
id. Тихая потеря событий из-за расхождения порядка присвоения и порядка коммита. - Запустить несколько relay без шардирования по агрегату. Получите переупорядочивание событий одной сущности и «обновление до создания» у потребителей.
- Верить в exactly-once на границе системы. EOS Kafka заканчивается там, где начинается ваш HTTP-вызов в платёжный шлюз.
- Генерировать
message_idв момент отправки. Ретрай породит новый id, дедупликация станет декорацией. Заодно не забывайте чиститьoutboxиinbox: обе таблицы растут монотонно и требуют партиционирования с retention больше горизонта ретраев. - Строить сагу там, где хватило бы одной транзакции. Разделение на два сервиса ради «микросервисности» превращает бесплатную атомарность в дорогую распределённую задачу. Про границы — ограниченные контексты.
Мини-итог
- Атомарность держится на одной записи коммита в одном журнале. Как только записей две, нужен протокол — и он всегда чем-то платит.
- Атомарный коммит не сводится к консенсусу: одно «нет» побеждает любое большинство, поэтому голос участника незаменим.
- 2PC блокирующий не по недосмотру, а по теореме Skeen–Stonebraker. Участник в
preparedобязан ждать — и это видно как зависшие строки вpg_prepared_xacts, неподвижныйoldest xminи растущий диск. - 3PC не решает проблему при разделении сети: он меняет зависание на нарушение атомарности. 2PC — это Paxos Commit с одним акцептором, и настоящее лечение — реплицировать решение консенсусом, как делают Spanner, TiDB, CockroachDB и Kafka.
- Сага меняет откат на компенсацию и изоляцию на её отсутствие. Порядок шагов (компенсируемые → pivot → повторяемые) — часть контракта. Компенсации нужно проектировать, тестировать и мониторить отдельно: их провал — штатный сценарий, а не исключение.
- Outbox сводит две границы атомарности к одной ценой гарантированных дубликатов. Дедупликация обязана коммититься вместе с эффектом.
- Exactly-once delivery невозможна (две генерала); exactly-once effect достижим через at-least-once плюс идемпотентность в общем атомарном шаге. Всё остальное — маркетинг с мелким шрифтом про границы системы.
Источники
- Jim Gray, Notes on Data Base Operating Systems (1978) — первоисточник по 2PC и проблеме двух генералов.
- Jim Gray, Leslie Lamport, Consensus on Transaction Commit (TODS, 2006) — 2PC как вырожденный Paxos Commit; ключевая статья темы.
- Dale Skeen, Michael Stonebraker, A Formal Model of Crash Recovery in a Distributed System (IEEE TSE, 1983) — доказательство блокирующего свойства; Skeen, Nonblocking Commit Protocols (SIGMOD 1981) — 3PC.
- Mohan, Lindsay, Obermarck, Transaction Management in the R* Distributed Database Management System (TODS, 1986) — presumed abort/commit.
- Hector Garcia-Molina, Kenneth Salem, Sagas (SIGMOD 1987) — оригинальная модель саги.
- Fischer, Lynch, Paterson, Impossibility of Distributed Consensus with One Faulty Process (JACM, 1985).
- Corbett et al., Spanner: Google’s Globally-Distributed Database (OSDI 2012); Peng, Dabek, Large-scale Incremental Processing Using Distributed Transactions and Notifications — Percolator (OSDI 2010).
- Pat Helland, Life Beyond Distributed Transactions: An Apostate’s Opinion (CIDR 2007) — почему в больших системах живут без распределённых транзакций.
- Martin Kleppmann, «Designing Data-Intensive Applications», глава 9 («Consistency and Consensus», раздел про атомарный коммит) и Is Kafka Exactly Once Really Exactly Once? — исходное описание EOS от Confluent.
- Документация: PostgreSQL
PREPARE TRANSACTION, Kafka transactions (KIP-98), Debezium Outbox Event Router, Flink end-to-end exactly-once.
Что дальше
Все три механизма этой статьи — 2PC с его ретраями COMMIT PREPARED, сага с повторами шагов и компенсаций, outbox с повторной публикацией — упираются в одно и то же требование: обработчик обязан выдерживать повтор. Дальше разбираем это требование системно: что такое at-most-once и at-least-once на уровне протокола, почему exactly-once — свойство эффекта, а не канала, как строить ключи идемпотентности, где хранить окно дедупликации и как не сломать его при масштабировании — Гарантии доставки и идемпотентность: at-least-once и миф exactly-once.