Transaction outbox: что он гарантирует, чего не гарантирует и во что обходится
Событие становится строкой той же базы, и тогда у него та же судьба, что у заказа. Дальше начинается интересное: релей, который помнит номер последней отправленной строки, теряет события молча, а обещанный порядок доставки паттерн не даёт — ни через опрос таблицы, ни через чтение WAL.
Полное техническое изложение
TL;DR
- Записать в базу и отправить событие — задача без решения в лоб: между двумя записями всегда есть момент, в который можно упасть.
- Outbox делает событие строкой той же базы. Один коммит — либо обе записи, либо ни одной.
- Стоит это +30 % к транзакции. Способ, которым outbox обычно избегают — отдельная транзакция для события, — стоит вдвое дороже и при этом небезопасен.
- Ошибка, которая не даёт о себе знать: релей, который помнит
last_id. Он теряет события молча. - Порядок паттерн не даёт, как только две транзакции пересеклись во времени: приходит порядок коммитов, а не вставок. Сохранить его можно только внутри одного ключа и только при одном релее.
- Дубли неизбежны — идемпотентность нужна потребителю.
Задача
db.insert(order)
db.commit()
broker.publish(OrderPlaced(order.id)) # <- здесь процесс умерЗаказ есть, события нет — потребитель о нём не узнает никогда. Поменяйте строки местами: событие есть, заказа нет — потребитель обработает призрак.
Третьего порядка не существует.
Решение
Событие становится строкой той же базы:
CREATE TABLE outbox (
id bigserial PRIMARY KEY,
payload jsonb NOT NULL,
published_at timestamptz
);db.insert(order)
db.insert(Outbox(payload={"type": "OrderPlaced", "order_id": order.id}))
db.commit() # обе строки или ни однойОтдельный процесс — релей — читает строки и отправляет их в брокер.
Что это стоит
Строка outbox добавляет к транзакции 30 %: она едет в тот же WAL и коммитится
тем же fsync. Отдельная транзакция для события стоит вдвое дороже — два
fsync вместо одного.
То есть небезопасный способ ещё и медленнее.
Ошибка, которая не даёт о себе знать
Самый естественный релей помнит, до какого номера дошёл:
SELECT id, payload FROM outbox WHERE id > :last_id ORDER BY id; -- неправильноbigserial выдаёт номер при INSERT, а видимой строка становится при
COMMIT. Транзакция, начавшаяся раньше, коммитится позже — и появляется
позади уже прочитанного. Отметка ушла вперёд, строка потеряна навсегда.
Ни исключения, ни лога. Метрика отставания показывает ноль.
Правильно — не помнить ничего:
SELECT id, payload FROM outbox
WHERE published_at IS NULL
ORDER BY id LIMIT 100
FOR UPDATE SKIP LOCKED;SKIP LOCKED заодно разводит два релея: первый берёт строки 1–3, второй 4–6,
пересечения нет.
Чего паттерн не даёт
Порядок. Пока транзакции не пересекаются во времени, порядок вставок и
порядок коммитов совпадают. Стоит двум пересечься — и всё: приложение вставило
A, потом B, а доставлены будут B, потом A. Через чтение WAL то же самое.
Сохранить порядок можно только внутри одного ключа: положите идентификатор
заказа в ключ сообщения, и события одного заказа придут по порядку — пока
релей один, потому что два релея со SKIP LOCKED обгоняют друг друга.
«Ровно один раз». Между отправкой и пометкой тоже можно упасть. Поменять их местами нельзя — тогда событие потеряется, а это хуже. Значит, дубли неизбежны, и потребитель должен быть идемпотентным:
INSERT INTO processed (event_id) VALUES (:id) ON CONFLICT DO NOTHING;
-- вставилось 0 строк -> дубль, работу не делаемВставка идентификатора и сама обработка — одной транзакцией. Иначе вы построили ту же двойную запись, с которой начали.
Два обязательных пункта эксплуатации
Частичный индекс. Релей почти всё время спрашивает «есть ли работа» и получает «нет». Без индекса на таблице в 71 МБ это 24 мс на каждый холостой опрос.
CREATE INDEX outbox_unpublished ON outbox (id) WHERE published_at IS NULL;53 микросекунды вместо 24 миллисекунд. И индекс не растёт вместе с таблицей: при пустой очереди — 8192 байта против 71 МБ у таблицы.
Уборка. DELETE места не возвращает: полмиллиона строк удаляются за
0,55 с, таблица как весила 71 МБ, так и весит. VACUUM вернёт страницы под
повторную запись (после него 11 МБ), но операционной системе отдаст только
хвост файла. При заметном потоке — секции по времени и DROP секции.
Когда не нужен
Если получатель — та же база, пишите в одной транзакции. Если событие можно потерять без последствий, отправляйте напрямую.
Обязателен он там, где по событию произойдёт необратимое действие — списание, отгрузка, письмо, — а расхождение обнаружится не сразу.
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. - «Ровно один раз» не бывает. Между отправкой и пометкой тоже можно упасть, поэтому идемпотентность нужна потребителю, а не отправителю.
Задача, у которой нет решения в лоб
Обработчик делает две вещи: пишет заказ в базу и сообщает о нём брокеру.
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.
Решение состоит в том, чтобы служба, отправляющая сообщение, сначала сохранила его в базе в рамках той же транзакции, которая обновляет бизнес-сущности. Отдельный процесс затем отправляет сообщения в брокер.
CREATE TABLE outbox (
id bigserial PRIMARY KEY,
payload jsonb NOT NULL,
published_at timestamptz
);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 поэтому шире трёх столбцов:
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.
если запись писателя содержит поле с именем, отсутствующим в записи читателя, значение писателя для этого поля игнорируется.
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 удалены и должны быть построены заново; произошло повреждение данных из-за ошибки конфигурации или иной проблемы.
И тут же — предложение, ради которого этот раздел написан:
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.
Сделайте потребителя идемпотентным, заставив его записывать идентификаторы обработанных сообщений в базу. При обработке сообщения потребитель может обнаружить и отбросить дубли, обратившись к базе.
И вот здесь id bigserial из нашего DDL перестаёт годиться. Дедупликация
работает по стабильному идентификатору события, а bigserial при
переигрывании выдаст новый номер — тот же INSERT … ON CONFLICT DO NOTHING
пропустит дубль как новое событие. Поэтому у Debezium id — uuid, задаваемый
при вставке, а не последовательностью. Менять это дешевле до первого
переигрывания, чем после.
Второе следствие, и оно противоречит разделу про рост таблицы ниже. Совет
удалять отправленные строки — DROP секции вместо DELETE — экономит место
ровно за счёт возможности переиграть. Удалённая секция — это удалённая история.
Решать здесь приходится один вопрос, и лучше явно: сколько времени назад вы
готовы переиграть события. Отсюда и срок хранения секций, а не наоборот.
Про сторону потребителя (сдвинуть смещение и прочитать топик заново) — то же самое с другой стороны: дубли создаёт не только повторная отправка, но и повторное чтение, и защищает от обеих одна и та же таблица обработанных идентификаторов.
Сколько это стоит
Возражение против outbox почти всегда одно: лишняя запись в горячем пути. Возражение проверяемое.
Строка outbox добавляет к транзакции 30 % — она едет в тот же WAL и
коммитится тем же fsync, что и заказ. А популярный способ «не грузить
транзакцию» — записать событие отдельно — стоит вдвое дороже, потому что
в нём два коммита и два fsync. То есть небезопасный вариант здесь ещё и
медленнее безопасного.
Это тот случай, когда интуиция подводит ровно наоборот: люди избегают outbox из соображений производительности и получают решение медленнее.
Релей: два способа, и между ними не «вкус»
Publish messages by polling the database's outbox table.
Публикуйте сообщения, опрашивая таблицу outbox.
Tail the database transaction log and publish each message/event inserted into
the outbox to the message broker.
Читайте журнал транзакций базы и публикуйте в брокер каждое сообщение, вставленное в outbox.
Разница между ними не в надёжности: доставят всё и тот и другой. Различий пять, и выбирают обычно по последнему — по цене эксплуатации:
| опрос таблицы | чтение журнала | |
|---|---|---|
| работает на | 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 работает как очередь; изменения записей в ней не допускаются.
То есть таблица под эти два способа выглядит по-разному, и переехать с одного на другой — не смена процесса, а миграция схемы.
Ошибка, которая стоит дороже всех
Самая естественная реализация опрашивающего релея — помнить, до какого номера дошёл:
SELECT id, payload FROM outbox WHERE id > :last_id ORDER BY id;Она неправильна, и ломается не в редком случае, а в обычном.
Причина в том, что два момента жизни строки разъехались. bigserial выдаёт
номер при INSERT — из последовательности, вне транзакции. Видимой для
других строка становится при COMMIT. Между этими двумя моментами проходит
столько времени, сколько работает транзакция, и за это время кто угодно
успевает получить номер побольше и закоммититься раньше.
Дальше арифметика: релей увидел строку 2, поднял отметку до 2, а строка 1
появилась после. Условие id > 2 её не найдёт уже никогда.
Ни исключения, ни записи в логе. Метрика «отставание релея», построенная как
max(id) - last_id, показывает ноль: с её точки зрения всё отправлено.
Лечится это отказом от памяти — состояние переносится в саму строку:
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.
Сообщения отправляются в брокер в том порядке, в каком их отправило приложение.
Обещание выполняется, пока транзакции не пересекаются во времени: тогда порядок вставок и порядок коммитов совпадают, и спорить не о чем. Стоит двум пересечься — и оно перестаёт быть верным. Опыт нужен ровно тот же, что в предыдущем разделе: приложение вставило 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.
Релей сообщений может опубликовать сообщение более одного раза.
Поменять их местами нельзя: тогда падение между ними потеряет событие, а это строго хуже дубля. Поэтому выбор сделан в пользу дубля, и он окончательный — на уровне отправителя эта задача не решается вовсе.
Значит, идемпотентность обязана быть у потребителя. Обычный способ — таблица обработанных идентификаторов:
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 выбранные строки, которые не удаётся заблокировать немедленно, пропускаются. Пропуск заблокированных строк даёт несогласованное представление данных, поэтому он не годится для работы общего назначения, но может использоваться, чтобы избежать состязания за блокировки при обращении нескольких потребителей к таблице-очереди.
Проверено: с 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.
Частичный индекс — это индекс, построенный по подмножеству таблицы; подмножество задаётся условным выражением, которое называют предикатом частичного индекса. Индекс содержит записи только для тех строк таблицы, которые удовлетворяют предикату.
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.
Потенциально чреват ошибками, поскольку разработчик может забыть опубликовать сообщение после обновления базы.
Недостаток назван точно: строку в 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 превращает сообщение в строку той же базы, и тогда у него та же судьба, что у заказа: один коммит — либо оба, либо ни одного.
- Это дёшево. Замерено на PostgreSQL 16.13: строка outbox добавляет к транзакции 30 %, а способ, которым её обычно избегают — записать событие отдельной транзакцией, — стоит вдвое дороже. Небезопасный способ здесь ещё и медленнее.
- Ошибка, которая не даёт о себе знать, — релей, который помнит номер последней отправленной строки. Он теряет события молча: без исключения, без записи в логе, при нулевом отставании в метрике.
- Порядок доставки паттерн не даёт — как только две транзакции пересеклись во времени. Обещание на канонической странице паттерна — Messages are sent to the message broker in the order they were sent by the application — выполняется, только пока они не пересекаются; стоит им пересечься, и приходит порядок коммитов. Ни через опрос таблицы, ни через чтение WAL.
- «Ровно один раз» не бывает. Между отправкой и пометкой тоже можно упасть, поэтому идемпотентность нужна потребителю, а не отправителю.
На самом деле
- Не при падении брокера, а при падении вашего процесса между двумя записями. Недоступность брокера решается ретраями и без всякого outbox. Задача, которую он решает, названа на канонической странице точно: How to atomically update the database and send messages to a message broker? Атомарность — вот слово, вокруг которого всё построено.
- Замерено на PostgreSQL 16.13: заказ без события — 301 мкс, тот же заказ со строкой outbox в одной транзакции — 392, то есть плюс 30 %. Строка едет в тот же WAL и коммитится тем же fsync. А способ, которым эту «лишнюю запись» обычно обходят — записать событие отдельной транзакцией, — стоит 632 мкс, вдвое дороже, потому что покупает второй fsync. Небезопасный вариант здесь ещё и медленнее безопасного.
- Эта ошибка не даёт ни исключения, ни лога, ни отставания в метрике — тем она и дорога.
bigserialвыдаёт номер приINSERT, а видимой строка становится приCOMMIT; транзакция, получившая номер раньше, может закоммититься позже и появиться позади уже прочитанного. Воспроизведено: отметка ушла на 2, строка с номером 1 не будет отправлена никогда. Ни исключения, ни лога, и метрикаmax(id) − last_idпоказывает ноль. - Каноническая страница паттерна обещает 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 против решения под конкретную базу), в форме таблицы (Debezium требует All changes in an outbox table are expected to be INSERT operations, то есть столбца состояния там быть не должно) и в эксплуатации: отставший слот репликации не даёт удалять WAL и может забить диск.
- Между отправкой в брокер и пометкой строки тоже можно упасть, и это признано в самом паттерне: The Message relay might publish a message more than once. Поменять две операции местами нельзя — тогда падение потеряет событие, что строго хуже. Выбор в пользу дубля окончательный, и идемпотентность нужна потребителю: вставка идентификатора события и сама обработка — одной транзакцией.
- Удалять нужно, но места это не возвращает. Замерено: полмиллиона отправленных строк удаляются за 0,55 с, а таблица как весила 71 МБ, так и весит — строки лишь помечены мёртвыми.
VACUUMвернёт страницы под повторную запись (0,10 с, после него 11 МБ), но операционной системе отдаст только хвост файла. При заметном потоке помогают секции по времени иDROPсекции.
Что разобрано
- Задача, у которой нет решения в лоб
- Решение: сделать событие строкой той же базы
- `payload jsonb` — это не схема, а её отсутствие
- Переиграть события: почему это не «просто прочитать таблицу снова»
- Сколько это стоит
- Релей: два способа, и между ними не «вкус»
- Ошибка, которая стоит дороже всех
- Чего паттерн не даёт: порядок
- Чего паттерн не даёт: «ровно один раз»
- Два релея
- Таблица, которая растёт
- Когда outbox не нужен
- Как воспроизвести числа
- Чем измерено
Расхожие заблуждения
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 секции.
Проверьте себя
Обработчик пишет заказ, коммитит, затем отправляет событие в брокер. Процесс умер между коммитом и отправкой. Что произошло?
Источники и что читать дальше
9 ИСТОЧНИКОВ
- 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
- 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
- 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
- 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
- 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
- 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
- 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/
- 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/
- 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