Deep Engineering
Продвинутый·Опубликовано·30 МИН

Transaction outbox: что он гарантирует, чего не гарантирует и во что обходится

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

Полное техническое изложение

TL;DR

Задача звучит невинно: записать заказ в базу и сообщить о нём остальным. Двух систем достаточно, чтобы это стало неразрешимо в лоб — между двумя записями всегда есть момент, в который можно упасть.

  • Outbox превращает сообщение в строку той же базы, и тогда у него та же судьба, что у заказа: один коммит — либо оба, либо ни одного.
  • Это дёшево. Замерено на PostgreSQL 16.13: строка outbox добавляет к транзакции 30 %, а способ, которым её обычно избегают — записать событие отдельной транзакцией, — стоит вдвое дороже. Небезопасный способ здесь ещё и медленнее.
  • Ошибка, которая не даёт о себе знать, — релей, который помнит номер последней отправленной строки. Он теряет события молча: без исключения, без записи в логе, при нулевом отставании в метрике.
  • Порядок доставки паттерн не даёт — как только две транзакции пересеклись во времени. Обещание на канонической странице паттерна — Messages are sent to the message broker in the order they were sent by the application (Сообщения отправляются в брокер в том порядке, в каком их отправило приложение) — выполняется, только пока они не пересекаются; стоит им пересечься, и приходит порядок коммитов. Ни через опрос таблицы, ни через чтение WAL.
  • «Ровно один раз» не бывает. Между отправкой и пометкой тоже можно упасть, поэтому идемпотентность нужна потребителю, а не отправителю.

Задача, у которой нет решения в лоб

Обработчик делает две вещи: пишет заказ в базу и сообщает о нём брокеру.

PYTHON
def place_order(request):
    order = Order(total=request.total)
    db.insert(order)
    db.commit()
    broker.publish(OrderPlaced(order.id))

Между commit() и publish() процесс может умереть. Не «теоретически»: деплой, OOM-killer, потеря сети до брокера, таймаут. Поменяйте строки местами — станет не лучше, а иначе.

Оба порядка неправильны, и по-разному. Потерянное событие тихо: обе системы считают, что у них всё в порядке, и узнаете вы об этом от пользователя. Событие-призрак шумно: потребитель спишет резерв или отправит письмо про заказ, которого нет.

Третьего порядка не существует. Пока брокер не участвует в транзакции базы, между двумя записями всегда есть момент.

Распределённая транзакция на два ресурса эту задачу решала бы — но платить за неё придётся координатором и блокировками на время согласования. А для брокера это и вовсе не вариант: Kafka в двухфазном коммите не участвует.

Решение: сделать событие строкой той же базы

The solution is for the service that sends the message to first store the message in the database as part of the transaction that updates the business entities. A separate process then sends the messages to the message broker.

перевод

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

Крис Ричардсон, Pattern: Transactional outbox
SQL
CREATE TABLE outbox (
    id           bigserial PRIMARY KEY,
    payload      jsonb       NOT NULL,
    published_at timestamptz
);
PYTHON
def place_order(request):
    order = Order(total=request.total)
    db.insert(order)
    db.insert(Outbox(payload={"type": "OrderPlaced", "order_id": order.id}))
    db.commit()          # обе строки или ни одной

Обратите внимание, чего здесь нет. Нет ничего специфического для паттерна: это обычная атомарность транзакции. Всё изобретение — в том, что сообщение перестало быть сообщением и стало строкой.

Дальше нужен второй процесс — релей, который читает строки и отправляет их в брокер. Именно в нём и живут все настоящие вопросы.

payload jsonb — это не схема, а её отсутствие

В таблице выше payload — столбец jsonb, и это упрощение, за которое платят позже. Релей его содержимого не читает и не проверяет:

The outbox event router SMT supports arbitrary payload formats. The SMT passes on payload column values that it reads from the outbox table without modification.

перевод

Преобразование поддерживает произвольные форматы полезной нагрузки. Оно передаёт значения столбца payload, прочитанные из таблицы outbox, без изменений.

Debezium — Outbox Event Router

То есть между отправителем и потребителем нет ни одной точки, где несовместимое изменение было бы замечено. Оно проявится у потребителя, в продакшене, на сообщениях, которые уже нельзя не отправить.

Ссылочная форма таблицы у Debezium поэтому шире трёх столбцов:

id            | uuid                   | not null
aggregatetype | character varying(255) | not null
aggregateid   | character varying(255) | not null
type          | character varying(255) | not null
payload       | jsonb                  |

Каждый столбец здесь отвечает за что-то, что иначе придётся решать в коде. aggregatetype определяет имя топика, aggregateid становится ключом сообщения (а значит — разделом и порядком), type уезжает заголовком, а id — тем, по чему потребитель убирает дубли: You can use this ID, for example, to remove duplicate messages (Этот идентификатор можно использовать, например, чтобы убирать дубли сообщений).

Разница с bigserial из нашего DDL здесь принципиальная, и к ней мы вернёмся в следующем разделе.

Отдельно про сам конверт. Как только полей становится больше одного, значение сообщения перестаёт быть голой нагрузкой:

A representation of the outbox change event. The default structure is JSON. By default, the Kafka message value is solely comprised of the payload value. However, if the outbox event is configured to include additional fields, the Kafka message value contains an envelope encapsulating both payload and the additional fields, and each field is represented separately.

перевод

Представление outbox-события. Структура по умолчанию — JSON. По умолчанию значение Kafka-сообщения состоит исключительно из значения payload. Однако если outbox-событие настроено на включение дополнительных полей, значение Kafka-сообщения содержит конверт, оборачивающий и полезную нагрузку, и дополнительные поля, причём каждое поле представлено отдельно.

Там же

Что делать со сменой схемы. Правила совместимости придумывать не надо — они записаны в форматах, которые для этого и сделаны. У Avro они сформулированы как правила чтения, и из них прямо следуют обе совместимости:

if the writer's record contains a field with a name not present in the reader's record, the writer's value for that field is ignored.

перевод

если запись писателя содержит поле с именем, отсутствующим в записи читателя, значение писателя для этого поля игнорируется.

Apache Avro 1.11.1 — Schema Resolution

if the reader's record schema has a field that contains a default value, and writer's schema does not have a field with the same name, then the reader should use the default value from its field.

перевод

если схема записи читателя имеет поле со значением по умолчанию, а схема писателя не содержит поля с тем же именем, читатель должен использовать значение по умолчанию из своего поля.

Там же

Первое правило — про старого читателя и нового писателя, второе — про нового читателя и старого писателя. Обоих не будет, если у нового поля нет значения по умолчанию: тогда старое сообщение просто не прочитается.

У протокол-буферов те же три правила названы прямым текстом: Adding new fields is safe (Добавлять новые поля безопасно), Removing fields is safe (Удалять поля безопасно) и, в противоположность им, Changing field numbers for any existing field is not safe (Менять номера у любого существующего поля небезопасно). Плюс отдельная оговорка, которой в JSON-схеме соответствия нет: удалённый номер нельзя переиспользовать, для этого есть reserved.

Практическое правило для outbox складывается из этого само: менять смысл существующего поля нельзя, добавлять поле со значением по умолчанию можно, удалять — можно, если ни один потребитель его не требует. Всё остальное — новый type, а не новая версия старого. И если нужен не договор на словах, а проверка, Debezium называет способ:

Using Avro can be beneficial for message format governance and for ensuring that outbox event schemas evolve in a backwards-compatible way.

перевод

Использование Avro может быть полезно для управления форматом сообщений и для того, чтобы схемы outbox-событий эволюционировали обратносовместимым образом.

Там же

Переиграть события: почему это не «просто прочитать таблицу снова»

Рано или поздно топик придётся наполнить заново: потребитель появился позже событий, топик удалили, данные испортили. Debezium перечисляет поводы прямо:

You might want to perform an ad hoc snapshot after any of the following changes occur in your Debezium environment: The connector configuration is modified to capture a different set of tables. Kafka topics are deleted and must be rebuilt. Data corruption occurs due to a configuration error or some other problem.

перевод

Такой снимок может понадобиться после любого из следующих изменений в вашей среде: конфигурация коннектора изменена так, чтобы захватывать другой набор таблиц; топики Kafka удалены и должны быть построены заново; произошло повреждение данных из-за ошибки конфигурации или иной проблемы.

Debezium — Ad hoc snapshots

И тут же — предложение, ради которого этот раздел написан:

When you initiate an ad hoc snapshot of an existing table, the connector appends content to the topic that already exists for the table.

перевод

Когда вы инициируете снимок по требованию для существующей таблицы, коннектор дописывает содержимое в уже существующий для этой таблицы топик.

Там же

Дописывает, а не заменяет. Потребитель получает вторую копию всего. Это не дефект инструмента: у топика нет операции «переписать», и быть её не может. Значит, вся тяжесть переигрывания ложится на дедупликацию у потребителя — ту самую, которая обязательна из-за доставки «не менее одного раза» — ей посвящён отдельный раздел ниже:

Make a consumer idempotent by having it record the IDs of processed messages in the database. When processing a message, a consumer can detect and discard duplicates by querying the database.

перевод

Сделайте потребителя идемпотентным, заставив его записывать идентификаторы обработанных сообщений в базу. При обработке сообщения потребитель может обнаружить и отбросить дубли, обратившись к базе.

Chris Richardson — Pattern: Idempotent Consumer

И вот здесь id bigserial из нашего DDL перестаёт годиться. Дедупликация работает по стабильному идентификатору события, а bigserial при переигрывании выдаст новый номер — тот же INSERT … ON CONFLICT DO NOTHING пропустит дубль как новое событие. Поэтому у Debezium iduuid, задаваемый при вставке, а не последовательностью. Менять это дешевле до первого переигрывания, чем после.

Второе следствие, и оно противоречит разделу про рост таблицы ниже. Совет удалять отправленные строки — DROP секции вместо DELETE — экономит место ровно за счёт возможности переиграть. Удалённая секция — это удалённая история. Решать здесь приходится один вопрос, и лучше явно: сколько времени назад вы готовы переиграть события. Отсюда и срок хранения секций, а не наоборот.

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

Сколько это стоит

Возражение против outbox почти всегда одно: лишняя запись в горячем пути. Возражение проверяемое.

Строка outbox добавляет к транзакции 30 % — она едет в тот же WAL и коммитится тем же fsync, что и заказ. А популярный способ «не грузить транзакцию» — записать событие отдельно — стоит вдвое дороже, потому что в нём два коммита и два fsync. То есть небезопасный вариант здесь ещё и медленнее безопасного.

Это тот случай, когда интуиция подводит ровно наоборот: люди избегают outbox из соображений производительности и получают решение медленнее.

Релей: два способа, и между ними не «вкус»

Publish messages by polling the database's outbox table.

перевод

Публикуйте сообщения, опрашивая таблицу outbox.

Крис Ричардсон, Pattern: Polling publisher

Tail the database transaction log and publish each message/event inserted into the outbox to the message broker.

перевод

Читайте журнал транзакций базы и публикуйте в брокер каждое сообщение, вставленное в outbox.

Крис Ричардсон, Pattern: Transaction log tailing

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

опрос таблицычтение журнала
работает наany SQL database (любой SQL-базе)конкретной базе, отдельно на каждую
задержкане меньше интервала опросаинтервала не ждёт: журнал толкает сам
холостой ходзапрос к базе каждый интервалничего: журнал сам толкает
форма таблицынужен столбец состояния и UPDATEтолько INSERT, состояния нет
эксплуатацияобычный процессслот репликации, который растит WAL, если отстанет

Debezium, основная реализация чтения журнала, требование к таблице формулирует прямо:

All changes in an outbox table are expected to be INSERT operations. That is, an outbox table functions as a queue; updates to records in an outbox table are not allowed.

перевод

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

Debezium, Outbox Event Router

То есть таблица под эти два способа выглядит по-разному, и переехать с одного на другой — не смена процесса, а миграция схемы.

Ошибка, которая стоит дороже всех

Самая естественная реализация опрашивающего релея — помнить, до какого номера дошёл:

SQL
SELECT id, payload FROM outbox WHERE id > :last_id ORDER BY id;

Она неправильна, и ломается не в редком случае, а в обычном.

Причина в том, что два момента жизни строки разъехались. bigserial выдаёт номер при INSERT — из последовательности, вне транзакции. Видимой для других строка становится при COMMIT. Между этими двумя моментами проходит столько времени, сколько работает транзакция, и за это время кто угодно успевает получить номер побольше и закоммититься раньше.

Дальше арифметика: релей увидел строку 2, поднял отметку до 2, а строка 1 появилась после. Условие id > 2 её не найдёт уже никогда.

Ни исключения, ни записи в логе. Метрика «отставание релея», построенная как max(id) - last_id, показывает ноль: с её точки зрения всё отправлено.

Лечится это отказом от памяти — состояние переносится в саму строку:

SQL
SELECT id, payload FROM outbox
  WHERE published_at IS NULL
  ORDER BY id
  LIMIT 100
  FOR UPDATE SKIP LOCKED;

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

Чего паттерн не даёт: порядок

На канонической странице паттерна в списке достоинств стоит:

Messages are sent to the message broker in the order they were sent by the application.

перевод

Сообщения отправляются в брокер в том порядке, в каком их отправило приложение.

Крис Ричардсон, Pattern: Transactional outbox

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

Через чтение журнала — то же самое. test_decoding показывает транзакцию со строкой id = 2 раньше транзакции со строкой id = 1, причём у той, что пришла первой, номер транзакции больше: она началась позже, а закоммитилась раньше. Значит, порядок и не по номеру строки, и не по номеру транзакции, а по моменту коммита.

Ту же трудность признаёт и страница про опрашивающий релей — в списке недостатков у неё стоит Tricky to publish events in order (Непросто публиковать события по порядку). Измерение уточняет её природу: это не трудность реализации, которую можно преодолеть аккуратностью, а свойство любого релея на PostgreSQL.

Что с этим делать. Общего порядка у вас не будет — на него нельзя опираться, и потребителя нельзя писать так, будто он есть. Порядок можно сохранить только внутри одного ключа, и только пока таблицу читает один релей: два релея со SKIP LOCKED из раздела ниже ломают и его. На этом построена промышленная реализация: Debezium кладёт идентификатор агрегата в ключ сообщения Kafka — This is important for maintaining correct order in Kafka partitions (Это важно для сохранения правильного порядка в разделах Kafka). События одного заказа попадают в один раздел и приходят по порядку; события разных заказов между собой не упорядочены — и не должны быть.

Чего паттерн не даёт: «ровно один раз»

Релей делает две вещи: отправляет и помечает. Между ними, как и в самом начале статьи, можно упасть.

The Message relay might publish a message more than once.

перевод

Релей сообщений может опубликовать сообщение более одного раза.

Крис Ричардсон, Pattern: Transactional outbox

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

Значит, идемпотентность обязана быть у потребителя. Обычный способ — таблица обработанных идентификаторов:

SQL
INSERT INTO processed (event_id) VALUES (:id) ON CONFLICT DO NOTHING;
-- 0 строк вставлено -> это дубль, работу не делаем

Ключевое здесь — что вставка идентификатора и сама обработка идут одной транзакцией. Иначе вы построили ту же двойную запись, с которой начали.

Два релея

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

With SKIP LOCKED, any selected rows that cannot be immediately locked are skipped. Skipping locked rows provides an inconsistent view of the data, so this is not suitable for general purpose work, but can be used to avoid lock contention with multiple consumers accessing a queue-like table.

перевод

При SKIP LOCKED выбранные строки, которые не удаётся заблокировать немедленно, пропускаются. Пропуск заблокированных строк даёт несогласованное представление данных, поэтому он не годится для работы общего назначения, но может использоваться, чтобы избежать состязания за блокировки при обращении нескольких потребителей к таблице-очереди.

PostgreSQL 16, The Locking Clause

Проверено: с SKIP LOCKED первый релей взял строки 1–3, второй 4–6, пересечение пусто. Без него второй релей на lock_timeout = 300ms падает, не взяв ни строки, — то есть не ускоряет, а простаивает.

Оговорка в документации про несогласованное представление здесь не мелким шрифтом. Именно она и означает, что порядок сломан ещё раз: второй релей обгоняет первый. Для очереди это допустимо ровно потому, что порядка у нас уже нет.

Прежде чем добавлять второй релей, стоит покрутить ручку подешевле. Партия по одной строке даёт 1 627 строк/с, партия по тысяче — 104 656: разница почти вся состоит из коммитов, и один процесс с большой партией обгоняет несколько процессов с маленькой. И порядок внутри ключа при одном релее ещё цел.

Таблица, которая растёт

Две вещи, которые ломают outbox в эксплуатации, обе про размер.

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

Спасает частичный индекс.

A partial index is an index built over a subset of a table; the subset is defined by a conditional expression (called the predicate of the partial index). The index contains entries only for those table rows that satisfy the predicate.

перевод

Частичный индекс — это индекс, построенный по подмножеству таблицы; подмножество задаётся условным выражением, которое называют предикатом частичного индекса. Индекс содержит записи только для тех строк таблицы, которые удовлетворяют предикату.

PostgreSQL 16, Partial Indexes
SQL
CREATE INDEX outbox_unpublished ON outbox (id) WHERE published_at IS NULL;

24 миллисекунды превращаются в 53 микросекунды. И главное свойство: индекс не растёт вместе с таблицей. При пустой очереди он занимает 8192 байта — одну страницу, минимум, — тогда как таблица рядом весит 71 МБ. Обычный индекс по published_at так не умеет: он хранит запись для каждой строки.

Вторая — уборка. Отправленные строки надо удалять, и DELETE места не возвращает: полмиллиона строк удаляются за 0,55 с, а таблица как весила 71 МБ, так и весит. VACUUM вернёт страницы под повторную запись (0,10 с, после него 11 МБ), но операционной системе отдаст только хвост файла.

Отсюда обычный совет для заметного потока: секции по времени и DROP секции вместо DELETE. DROP удаляет файл целиком — и это единственная операция, которая возвращает место немедленно.

Когда outbox не нужен

Паттерн решает одну задачу: согласовать запись в базу с сообщением наружу. Если такой задачи нет, он лишний.

  • Если получатель — та же база, никакого outbox не нужно: пишите в одной транзакции.
  • Если сообщение можно потерять без последствий (метрика, лог, прогрев кеша), цена паттерна не окупается: отправляйте напрямую.
  • Если потребитель и так регулярно перечитывает состояние из базы, событие — всего лишь подсказка «загляни пораньше», и его потеря стоит задержки, а не расхождения.

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

Potentially error prone since the developer might forget to publish the message/event after updating the database.

перевод

Потенциально чреват ошибками, поскольку разработчик может забыть опубликовать сообщение после обновления базы.

Крис Ричардсон, Pattern: Transactional outbox

Недостаток назван точно: строку в outbox тоже можно забыть вставить. Разница в том, что забытая строка — ошибка одного места в коде, которую видно на ревью, а разъехавшаяся двойная запись — свойство архитектуры, которое не видно нигде.

Как воспроизвести числа

Два скрипта, оба на живой базе, оба печатают то, что вернул PostgreSQL.

createdb -p 5433 outbox_bench
python3 relay.py     # восемь наблюдений, ни одного замера времени
python3 cost.py      # четыре блока замеров

Скрипты лежат в каталоге замеров репозитория, рядом с описанием того, что получилось. Нужен psycopg 3; адрес базы берётся из OUTBOX_DSN. Пятый раздел relay.py читает WAL и требует wal_level = logical — при другом уровне он скажет об этом и пропустится, остальные семь работают на любом.

Опубликованный прогон: PostgreSQL 16.13, synchronous_commit = on, fsync = on, psycopg 3.3.4, Python 3.11.15, август 2026 года. Всё измеренное упирается в fsync, поэтому на другом диске абсолютные значения будут другими — переносятся отношения внутри блока и порядок величин.

Чем измерено

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

Расхожие заблуждения

Утверждение

Outbox нужен, чтобы не потерять сообщение при падении брокера

На самом деле

Не при падении брокера, а при падении вашего процесса между двумя записями. Недоступность брокера решается ретраями и без всякого outbox. Задача, которую он решает, названа на канонической странице точно: How to atomically update the database and send messages to a message broker? (Как атомарно обновить базу данных и отправить сообщения в брокер?) Атомарность — вот слово, вокруг которого всё построено.

Утверждение

Строка в outbox — это лишняя запись, которая замедляет горячий путь

На самом деле

Замерено на PostgreSQL 16.13: заказ без события — 301 мкс, тот же заказ со строкой outbox в одной транзакции — 392, то есть плюс 30 %. Строка едет в тот же WAL и коммитится тем же fsync. А способ, которым эту «лишнюю запись» обычно обходят — записать событие отдельной транзакцией, — стоит 632 мкс, вдвое дороже, потому что покупает второй fsync. Небезопасный вариант здесь ещё и медленнее безопасного.

Утверждение

Релей достаточно написать как WHERE id > :last_id

На самом деле

Эта ошибка не даёт ни исключения, ни лога, ни отставания в метрике — тем она и дорога. bigserial выдаёт номер при INSERT, а видимой строка становится при COMMIT; транзакция, получившая номер раньше, может закоммититься позже и появиться позади уже прочитанного. Воспроизведено: отметка ушла на 2, строка с номером 1 не будет отправлена никогда. Ни исключения, ни лога, и метрика max(id) − last_id показывает ноль.

Утверждение

Outbox сохраняет порядок событий

На самом деле

Каноническая страница паттерна обещает Messages are sent to the message broker in the order they were sent by the application (Сообщения отправляются в брокер в том порядке, в каком их отправило приложение), и это выполняется, пока транзакции не пересекаются во времени. Стоит двум пересечься — и всё: приложение вставило A, потом B, а доставлены будут B, потом A, потому что B закоммитилась первой. Через чтение WAL то же самое, причём у той транзакции, что пришла первой, номер БОЛЬШЕ: она началась позже, а закоммитилась раньше. Ту же трудность признаёт и страница про опрашивающий релей — Tricky to publish events in order (Непросто публиковать события по порядку); измерение уточняет, что это не трудность реализации, а свойство любого релея на PostgreSQL. Порядок сохраняется только внутри одного ключа и только пока релей один.

Утверждение

Чтение журнала транзакций надёжнее опроса таблицы

На самом деле

Доедет всё и там, и там — разница не в надёжности. Она в задержке (опрос не может ответить быстрее своего интервала, журнал интервала не ждёт), в переносимости (any SQL database (любая SQL-база) против решения под конкретную базу), в форме таблицы (Debezium требует All changes in an outbox table are expected to be INSERT operations (все изменения в таблице outbox — это операции INSERT), то есть столбца состояния там быть не должно) и в эксплуатации: отставший слот репликации не даёт удалять WAL и может забить диск.

Утверждение

С outbox сообщение доставляется ровно один раз

На самом деле

Между отправкой в брокер и пометкой строки тоже можно упасть, и это признано в самом паттерне: The Message relay might publish a message more than once (Релей сообщений может опубликовать сообщение более одного раза). Поменять две операции местами нельзя — тогда падение потеряет событие, что строго хуже. Выбор в пользу дубля окончательный, и идемпотентность нужна потребителю: вставка идентификатора события и сама обработка — одной транзакцией.

Утверждение

Отправленные строки можно просто удалять запросом DELETE

На самом деле

Удалять нужно, но места это не возвращает. Замерено: полмиллиона отправленных строк удаляются за 0,55 с, а таблица как весила 71 МБ, так и весит — строки лишь помечены мёртвыми. VACUUM вернёт страницы под повторную запись (0,10 с, после него 11 МБ), но операционной системе отдаст только хвост файла. При заметном потоке помогают секции по времени и DROP секции.

Проверьте себя

Вопрос 1 из 5

Обработчик пишет заказ, коммитит, затем отправляет событие в брокер. Процесс умер между коммитом и отправкой. Что произошло?

Источники и что читать дальше

9 ИСТОЧНИКОВ

  1. Chris Richardson — Pattern: Transactional outboxИсточник. Каноническое описание паттерна. Задача: «How to atomically update the database and send messages to a message broker?» (Как атомарно обновить базу данных и отправить сообщения в брокер?). Решение: «The solution is for the service that sends the message to first store the message in the database as part of the transaction that updates the business entities. A separate process then sends the messages to the message broker» (Решение состоит в том, чтобы служба, отправляющая сообщение, сначала сохранила его в базе в рамках той же транзакции, которая обновляет бизнес-сущности. Отдельный процесс затем отправляет сообщения в брокер). В списке достоинств стоит и утверждение про порядок: «Messages are sent to the message broker in the order they were sent by the application» (Сообщения отправляются в брокер в том порядке, в каком их отправило приложение) — оно и проверяется в статье. В недостатках стоит ровно одна строка — «Potentially error prone since the developer might forget to publish the message/event after updating the database» (Потенциально чреват ошибками, поскольку разработчик может забыть опубликовать сообщение после обновления базы), она цитируется в финале статьи. Возможность дубля признана отдельным разделом Issues: «The Message relay might publish a message more than once» (Релей сообщений может опубликовать сообщение более одного раза).https://microservices.io/patterns/data/transactional-outbox.html
  2. Chris Richardson — Pattern: Polling publisherИсточник. Способ доставки через опрос таблицы: «Publish messages by polling the database's outbox table» (Публикуйте сообщения, опрашивая таблицу outbox). Достоинство одно: «Works with any SQL database» (Работает с любой SQL-базой). А в недостатках стоит ровно то, что противоречит обещанию порядка на странице самого паттерна: «Tricky to publish events in order» (Непросто публиковать события по порядку).https://microservices.io/patterns/data/polling-publisher.html
  3. Chris Richardson — Pattern: Transaction log tailingИсточник. Второй способ доставки: «Tail the database transaction log and publish each message/event inserted into the outbox to the message broker» (Читайте журнал транзакций базы и публикуйте в брокер каждое сообщение, вставленное в outbox). Из достоинств названы «No 2PC» (Без двухфазного коммита) и «Guaranteed to be accurate» (Гарантированно точен). Недостатков три, и для статьи важны два: «Requires database specific solutions» (Требует решений, специфичных для базы) и «Tricky to avoid duplicate publishing» (Непросто избежать повторной публикации); третий — «Relatively obscure although becoming increasing common» (Относительно малоизвестен, хотя встречается всё чаще).https://microservices.io/patterns/data/transaction-log-tailing.html
  4. PostgreSQL 16 — SELECT, The Locking ClauseОфициальная документация. Разрешение на то, чтобы два релея работали над одной таблицей: «With SKIP LOCKED, any selected rows that cannot be immediately locked are skipped. Skipping locked rows provides an inconsistent view of the data, so this is not suitable for general purpose work, but can be used to avoid lock contention with multiple consumers accessing a queue-like table» (При SKIP LOCKED выбранные строки, которые не удаётся заблокировать немедленно, пропускаются. Пропуск заблокированных строк даёт несогласованное представление данных, поэтому он не годится для работы общего назначения, но может использоваться, чтобы избежать состязания за блокировки при обращении нескольких потребителей к таблице-очереди). Оговорка про несогласованное представление здесь не мелким шрифтом: она и есть причина, по которой этот приём годится только для очереди.https://www.postgresql.org/docs/16/sql-select.html
  5. PostgreSQL 16 — Partial IndexesОфициальная документация. Определение: «A partial index is an index built over a subset of a table; the subset is defined by a conditional expression (called the predicate of the partial index). The index contains entries only for those table rows that satisfy the predicate» (Частичный индекс — это индекс, построенный по подмножеству таблицы; подмножество задаётся условным выражением, которое называют предикатом частичного индекса. Индекс содержит записи только для тех строк таблицы, которые удовлетворяют предикату). И причина, по которой он подходит очереди: «One major reason for using a partial index is to avoid indexing common values» (Одна из главных причин использовать частичный индекс — не индексировать часто встречающиеся значения).https://www.postgresql.org/docs/16/indexes-partial.html
  6. Debezium — Outbox Event RouterОфициальная документация. Реализация чтения журнала для outbox. Требование к форме таблицы, которое отличает этот путь от опроса: «All changes in an outbox table are expected to be INSERT operations. That is, an outbox table functions as a queue; updates to records in an outbox table are not allowed. The SMT automatically filters out DELETE operations on an outbox table» (Ожидается, что все изменения в таблице outbox — это операции INSERT. То есть таблица outbox работает как очередь; изменения записей в ней не допускаются. Преобразование автоматически отфильтровывает операции DELETE над таблицей outbox). И то, зачем в таблице отдельный столбец с идентификатором агрегата: «The SMT uses this value as the key in the emitted outbox message. This is important for maintaining correct order in Kafka partitions» (Преобразование использует это значение как ключ отправляемого сообщения. Это важно для сохранения правильного порядка в разделах Kafka).https://debezium.io/documentation/reference/stable/transformations/outbox-event-router.html
  7. Apache Avro 1.11.1 — Schema ResolutionОфициальная документация. Правила совместимости, записанные не как имена режимов, а как правила чтения. Вперёд: «if the writer's record contains a field with a name not present in the reader's record, the writer's value for that field is ignored» (если запись писателя содержит поле с именем, отсутствующим в записи читателя, значение писателя для этого поля игнорируется). Назад: «if the reader's record schema has a field that contains a default value, and writer's schema does not have a field with the same name, then the reader should use the default value from its field» (если схема записи читателя имеет поле со значением по умолчанию, а схема писателя не содержит поля с тем же именем, читатель должен использовать значение по умолчанию из своего поля). Обе перестают работать, если у нового поля нет значения по умолчанию.https://avro.apache.org/docs/1.11.1/specification/
  8. Protocol Buffers — Language Guide (proto 3), Updating A Message TypeОфициальная документация. Те же три правила, названные прямым текстом: «Adding new fields is safe» (Добавлять новые поля безопасно), «Removing fields is safe» (Удалять поля безопасно) и «Changing field numbers for any existing field is not safe» (Менять номера у любого существующего поля небезопасно). Плюс требование, которому в JSON-схеме нет соответствия: удалённый номер поля нельзя переиспользовать — для этого есть `reserved`.https://protobuf.dev/programming-guides/proto3/
  9. Chris Richardson — Pattern: Idempotent ConsumerИсточник. Способ дедупликации, к которому сводится вся защита от повторной отправки и повторного чтения: «Make a consumer idempotent by having it record the IDs of processed messages in the database» (Сделайте потребителя идемпотентным, заставив его записывать идентификаторы обработанных сообщений в базу). Отсюда требование к идентификатору события: он должен быть стабильным между прогонами, иначе при переигрывании дубль пройдёт как новое событие.https://microservices.io/patterns/communication-style/idempotent-consumer.html