MEASUREMENT
bench/outbox/relay.py
The script that produced the numbers in the article, and the record of the run. The file is read from the repository at build time — this is the code that was run, not a copy of it.
- Cited in
- /en/system-design/data-consistency/transaction-outbox
- Run on
- PostgreSQL 16.13, psycopg 3.3.4, Python 3.11.15
- How to run it
createdb -p 5433 outbox_bench python3 bench/outbox/relay.py python3 bench/outbox/cost.py
The run below is recorded in Russian. It is a lab record, kept in the language it was written in; the numbers, the tables and the code read the same either way.
Record of the run
Замеры для статьи «Transaction outbox: что он даёт и чего не даёт»
| Скрипт | Что делает |
|---|---|
relay.py |
восемь наблюдений без единого замера времени: чем плоха двойная запись, какую строку теряет релей с памятью, в каком порядке приходят события у опроса и у чтения WAL, как ведут себя два релея сразу, откуда берутся дубли |
cost.py |
цена: строка outbox в бизнес-транзакции, пропускная способность релея, холостой опрос, уборка |
createdb -p 5433 outbox_bench
python3 bench/outbox/relay.py
python3 bench/outbox/cost.py
Нужен psycopg 3 (pip install "psycopg[binary]"). DSN берётся из
OUTBOX_DSN, по умолчанию — локальный сокет на порту 5433. Раздел 5
relay.py читает WAL и требует wal_level = logical; при другом уровне он
честно скажет об этом и пропустится, остальные семь работают на любом.
Что здесь важно прочитать правильно
relay.py и cost.py отвечают на разные вопросы, и главный — в
relay.py. Спор про outbox почти никогда не про скорость: он про то,
какая строка потерялась и в каком порядке пришли события. Это проверяется
сравнением множеств, а не секундомером, поэтому замеров там нет вовсе.
По времени сопоставляются только три строки первого блока cost.py. Они
делают одну и ту же работу — записывают заказ и событие о нём — и меряются в
одном прогоне, чередуясь. Блоки 2–4 с ними и друг с другом не сопоставляются:
там другая работа и другие единицы.
Все числа упираются в fsync при коммите. На машине с другим диском
абсолютные значения будут другими; переносятся отношения внутри блока 1 и
порядок величин в блоках 2–4. Условия прогона печатаются в шапке скрипта:
synchronous_commit = on, fsync = on.
Что получилось (PostgreSQL 16.13, psycopg 3.3.4, Python 3.11.15)
Один заказ, медиана и 90-й перцентиль в микросекундах:
| как | мкс | p90 | отношение |
|---|---|---|---|
| заказ без события | 301 | 431 | ×1,00 |
| заказ + строка outbox, одна транзакция | 392 | 540 | ×1,30 |
| заказ и событие двумя транзакциями | 632 | 883 | ×2,10 |
Это главное число статьи: небезопасный способ здесь ещё и дороже безопасного, и ровно во столько раз, во сколько в нём больше коммитов.
Релей, 20 000 строк, отправки в брокер нет — только работа с базой:
| размер партии | строк/с | мкс на строку |
|---|---|---|
| 1 | 1 627 | 614,7 |
| 10 | 14 602 | 68,5 |
| 100 | 44 583 | 22,4 |
| 1 000 | 104 656 | 9,6 |
| 5 000 | 122 740 | 8,1 |
Холостой опрос — то, чем релей занят почти всё время. Таблица: 500 000 отправленных строк, 71 МБ, очередь пуста:
| индекс | медиана |
|---|---|
| без частичного индекса | 23 837 мкс |
с частичным индексом WHERE published_at IS NULL |
53 мкс (×447) |
Сам частичный индекс при пустой очереди занимает 8192 байта — минимальный размер, одна страница. Он хранит записи только для неотправленных строк, поэтому не растёт вместе с таблицей.
Уборка тех же 500 000 строк: DELETE — 0,55 с, размер таблицы после него
прежний, 71 МБ; VACUUM — 0,10 с, после него 11 МБ.
От прогона к прогону: первый блок 301–348 / 392–454 / 632–810 мкс (отношения устойчивы: ×1,30–1,31 и ×2,10–2,34). Второй блок сильно зависит от того, что ещё занимает диск: партия по одной строке давала от 442 до 1 627 строк/с. Опубликованный кадр снят на незанятой машине, все четыре блока — в одном прогоне.
Наблюдения relay.py
- Двойная запись ломается в обе стороны. «Коммит, потом отправка» при падении между ними оставляет заказ без события. Обратный порядок оставляет событие без заказа. Третьего порядка нет.
- Релей, который помнит
last_id, теряет события молча. Транзакция A получилаid = 1и ещё не закоммитилась; B получилаid = 2и закоммитилась первой. Релей увидел только строку 2, запомнил отметку 2 — и строка 1 не будет отправлена никогда. Ошибки нет, метрика отставания показывает ноль. Причина:bigserialвыдаёт номер приINSERT, а видимой строка становится приCOMMIT. - Столбец
published_atэту дыру закрывает: условие «не отправлено» находит строку заново, сколько бы она ни просидела незакоммиченной. - Но порядок доставки — порядок коммитов, а не номеров. В том же опыте
потребитель получает
['B', 'A']приid1 и 2. Приложение вставляло A первой. - Чтение WAL даёт ровно тот же порядок.
test_decodingпоказывает транзакцию со строкойid = 2РАНЬШЕ, чем транзакцию со строкойid = 1, причём у пришедшей первой номер транзакции БОЛЬШЕ — она началась позже, а закоммитилась раньше. Порядок не поidи не по номеру транзакции, а по моменту коммита. Разница между опросом и CDC — в цене и задержке, а не в том, что доедет. FOR UPDATE SKIP LOCKEDразводит два релея: первый взял[1, 2, 3], второй[4, 5, 6], пересечение пусто. БезSKIP LOCKEDвторой релей просто ждёт — наlock_timeout = 300msон падает, не взяв ни строки.- «Ровно один раз» не бывает. Падение между отправкой и пометкой даёт дубль; поменять их местами нельзя — тогда падение теряет событие. Outbox даёт «хотя бы один раз», и идемпотентность нужна потребителю.
- Незакоммиченная строка не загораживает более поздние. Событие,
закоммиченное вторым, релей видит сразу. Именно поэтому строка с меньшим
idпоявляется ПОЗАДИ уже прочитанной — и ломает релей с памятью.
Что не подтвердилось
Каноническая страница паттерна у Криса Ричардсона в списке достоинств
утверждает: «Messages are sent to the message broker in the order they were
sent by the application». Обещание выполняется, пока транзакции не
пересекаются во времени: тогда порядок вставок и порядок коммитов совпадают.
Стоит двум пересечься — и оно перестаёт быть верным: раздел 4 relay.py
показывает обратный порядок, раздел 5 — тот же обратный порядок в WAL. Ту же
трудность признаёт страница polling-publisher, где в недостатках стоит
«Tricky to publish events in order»; измерение уточняет, что это не трудность
реализации, а свойство любого релея на PostgreSQL. Порядок сохраняется только
внутри одного ключа и только пока таблицу читает один релей — раздел 6
показывает, как два релея со SKIP LOCKED обгоняют друг друга. На порядке
внутри ключа и построен Debezium, который кладёт aggregateid в ключ
сообщения Kafka.
Источники
- Chris Richardson. Pattern: Transactional outbox — https://microservices.io/patterns/data/transactional-outbox.html
- Chris Richardson. Pattern: Polling publisher — https://microservices.io/patterns/data/polling-publisher.html
- Chris Richardson. Pattern: Transaction log tailing — https://microservices.io/patterns/data/transaction-log-tailing.html
- PostgreSQL 16. SELECT, The Locking Clause — https://www.postgresql.org/docs/16/sql-select.html
- PostgreSQL 16. Partial Indexes — https://www.postgresql.org/docs/16/indexes-partial.html
- Debezium. Outbox Event Router — https://debezium.io/documentation/reference/stable/transformations/outbox-event-router.html
Script
404 lines"""Transaction outbox: что именно ломается и что именно чинится.
Здесь нет ни одного измерения времени — только наблюдения над настоящим
PostgreSQL. Каждый раздел печатает то, что вернула база, а не то, что про неё
принято рассказывать.
ЗАЧЕМ ОТДЕЛЬНЫЙ СКРИПТ БЕЗ ВРЕМЁН. Главные утверждения про outbox — не про
скорость, а про то, какая строка потерялась и в каком порядке пришли события.
Это проверяется сравнением множеств, и смешивать такую проверку с замерами
нельзя: время шумит, а «строка A не пришла» либо правда, либо нет.
ЗАПУСК:
createdb -p 5433 outbox_bench
python3 bench/outbox/relay.py
Нужен PostgreSQL с wal_level = logical — раздел 5 читает WAL. Если уровень
другой, раздел честно скажет об этом и пропустится; остальные семь работают
на любом уровне.
Снято на PostgreSQL 16.13, psycopg 3.3.4, Python 3.11.15.
"""
from __future__ import annotations
import os
import sys
import time
try:
import psycopg
except ImportError: # pragma: no cover
sys.exit("нужен psycopg 3: pip install 'psycopg[binary]'")
DSN = os.environ.get("OUTBOX_DSN", "postgres://postgres@/outbox_bench?host=/tmp&port=5433")
def head(n: int, title: str) -> None:
print(f"\n{'=' * 72}\n{n}. {title}\n{'=' * 72}")
def fresh(c: psycopg.Connection) -> None:
"""Пустые orders и outbox. Частичный индекс — см. cost.py, раздел 3."""
c.execute("DROP TABLE IF EXISTS outbox, orders")
c.execute("CREATE TABLE orders (id bigserial PRIMARY KEY, total int)")
c.execute(
"CREATE TABLE outbox ("
" id bigserial PRIMARY KEY,"
" payload jsonb NOT NULL,"
" published_at timestamptz)"
)
c.execute("CREATE INDEX outbox_unpub ON outbox (id) WHERE published_at IS NULL")
# ---------------------------------------------------------------------------
# 1. Двойная запись: почему без outbox не обойтись
# ---------------------------------------------------------------------------
def section_dual_write(c: psycopg.Connection) -> None:
head(1, "Двойная запись: два порядка, и оба неправильные")
print(
" Задача: записать заказ в базу И отправить событие в брокер.\n"
" Две системы, одна из них не участвует в транзакции базы.\n"
)
fresh(c)
c.commit()
broker: list[str] = []
# Порядок «сначала commit, потом publish». Падение между ними.
cur = c.cursor()
cur.execute("INSERT INTO orders (total) VALUES (4200)")
c.commit()
crashed_here = True # процесс умер ровно здесь
if not crashed_here: # pragma: no cover - ветка ради симметрии кода
broker.append("OrderPlaced")
rows = c.execute("SELECT count(*) FROM orders").fetchone()[0]
print(f" порядок «commit, потом publish», падение между ними:")
print(f" заказов в базе: {rows} событий в брокере: {len(broker)}")
print(" -> заказ есть, события нет. Потребитель о нём не узнает никогда.")
# Обратный порядок. Падение после отправки, до коммита.
fresh(c)
c.commit()
broker = []
cur.execute("INSERT INTO orders (total) VALUES (4200)")
broker.append("OrderPlaced")
c.rollback() # падение до коммита == откат
rows = c.execute("SELECT count(*) FROM orders").fetchone()[0]
print(f"\n порядок «publish, потом commit», падение между ними:")
print(f" заказов в базе: {rows} событий в брокере: {len(broker)}")
print(" -> событие о заказе есть, заказа нет. Потребитель обработает призрак.")
print(
"\n Третьего порядка не существует: пока брокер не участвует в\n"
" транзакции базы, между двумя записями всегда есть момент, в\n"
" который можно упасть."
)
# ---------------------------------------------------------------------------
# 2. Одна транзакция: обе записи или ни одной
# ---------------------------------------------------------------------------
def section_atomic(c: psycopg.Connection) -> None:
head(2, "Outbox: обе записи в одной транзакции")
fresh(c)
c.commit()
cur = c.cursor()
cur.execute("INSERT INTO orders (total) VALUES (4200)")
cur.execute("INSERT INTO outbox (payload) VALUES (%s::jsonb)", ('{"type":"OrderPlaced"}',))
c.rollback() # падение до коммита
o = c.execute("SELECT count(*) FROM orders").fetchone()[0]
b = c.execute("SELECT count(*) FROM outbox").fetchone()[0]
print(f" падение до коммита: заказов {o}, событий в outbox {b}")
cur.execute("INSERT INTO orders (total) VALUES (4200)")
cur.execute("INSERT INTO outbox (payload) VALUES (%s::jsonb)", ('{"type":"OrderPlaced"}',))
c.commit()
o = c.execute("SELECT count(*) FROM orders").fetchone()[0]
b = c.execute("SELECT count(*) FROM outbox").fetchone()[0]
print(f" падение после коммита: заказов {o}, событий в outbox {b}")
print(
"\n Расхождения нет ни в одном случае, и это не заслуга outbox как\n"
" паттерна — это обычная атомарность транзакции. Событие стало\n"
" строкой той же базы, поэтому у него та же судьба, что у заказа."
)
# ---------------------------------------------------------------------------
# 3. Релей по высшей отметке теряет строки
# ---------------------------------------------------------------------------
def section_high_water_mark(c: psycopg.Connection) -> None:
head(3, "Релей, который помнит last_id, теряет события молча")
fresh(c)
c.commit()
a = psycopg.connect(DSN)
b = psycopg.connect(DSN)
# A начал транзакцию и получил id=1, но ещё не закоммитился.
a.execute("INSERT INTO outbox (payload) VALUES ('{\"n\":\"A\"}'::jsonb)")
# B начал позже, получил id=2 — и закоммитился ПЕРВЫМ.
b.execute("INSERT INTO outbox (payload) VALUES ('{\"n\":\"B\"}'::jsonb)")
b.commit()
with psycopg.connect(DSN, autocommit=True) as r:
seen = r.execute("SELECT id, payload->>'n' FROM outbox WHERE id > 0 ORDER BY id").fetchall()
print(f" проход 1, релей видит: {seen}")
last = max(i for i, _ in seen)
print(f" релей запоминает last_id = {last}")
a.commit()
print(" теперь коммитится A — его строка имеет id = 1")
with psycopg.connect(DSN, autocommit=True) as r:
seen2 = r.execute(
"SELECT id, payload->>'n' FROM outbox WHERE id > %s ORDER BY id", (last,)
).fetchall()
table = r.execute("SELECT id, payload->>'n' FROM outbox ORDER BY id").fetchall()
print(f" проход 2 (id > {last}), релей видит: {seen2}")
print(f" а в таблице лежит: {table}")
print(
"\n Событие A не будет отправлено никогда. Ошибки нет, лога нет,\n"
" метрика «отставание релея» показывает ноль. Причина: bigserial\n"
" выдаёт номер при INSERT, а видимой строка становится при COMMIT,\n"
" и порядок этих двух событий не совпадает."
)
a.close()
b.close()
# ---------------------------------------------------------------------------
# 4. Статус вместо отметки — и что он не чинит
# ---------------------------------------------------------------------------
def section_status_column(c: psycopg.Connection) -> None:
head(4, "Столбец published_at: строка не теряется, но порядок — не по id")
fresh(c)
c.commit()
a = psycopg.connect(DSN)
b = psycopg.connect(DSN)
a.execute("INSERT INTO outbox (payload) VALUES ('{\"n\":\"A\"}'::jsonb)")
b.execute("INSERT INTO outbox (payload) VALUES ('{\"n\":\"B\"}'::jsonb)")
b.commit()
delivered: list[str] = []
with psycopg.connect(DSN, autocommit=True) as r:
rows = r.execute(
"SELECT id, payload->>'n' FROM outbox WHERE published_at IS NULL"
" ORDER BY id FOR UPDATE SKIP LOCKED"
).fetchall()
delivered += [n for _, n in rows]
r.execute(
"UPDATE outbox SET published_at = now() WHERE id = ANY(%s)", ([i for i, _ in rows],)
)
print(f" проход 1 отправил: {[n for _, n in rows]}")
a.commit()
with psycopg.connect(DSN, autocommit=True) as r:
rows2 = r.execute(
"SELECT id, payload->>'n' FROM outbox WHERE published_at IS NULL ORDER BY id"
).fetchall()
delivered += [n for _, n in rows2]
print(f" проход 2 отправил: {[n for _, n in rows2]}")
print(f"\n потребитель получил: {delivered}")
print(
" ничего не потеряно — условие «не отправлено» ищет строку заново,\n"
" сколько бы времени она ни просидела незакоммиченной.\n"
)
print(
f" Но порядок доставки {delivered}, а id у них 1 и 2. То есть\n"
" outbox выдаёт события в порядке КОММИТОВ, а не в порядке номеров.\n"
" Общего порядка он не сохраняет — сохранить его можно только\n"
" внутри одного ключа, разложив потребителей по ключу."
)
a.close()
b.close()
# ---------------------------------------------------------------------------
# 5. Тот же опыт, но через WAL
# ---------------------------------------------------------------------------
def section_logical_decoding(c: psycopg.Connection) -> None:
head(5, "Чтение WAL (CDC): тот же порядок коммитов, без опроса таблицы")
level = c.execute("SHOW wal_level").fetchone()[0]
if level != "logical":
print(f" wal_level = {level}, а нужен logical — раздел пропущен.")
print(" Включается в postgresql.conf и требует перезапуска сервера.")
return
fresh(c)
c.commit()
with psycopg.connect(DSN, autocommit=True) as s:
exists = s.execute(
"SELECT 1 FROM pg_replication_slots WHERE slot_name = 'outbox_bench_slot'"
).fetchone()
if exists:
s.execute("SELECT pg_drop_replication_slot('outbox_bench_slot')")
s.execute("SELECT pg_create_logical_replication_slot('outbox_bench_slot','test_decoding')")
a = psycopg.connect(DSN)
b = psycopg.connect(DSN)
a.execute("INSERT INTO outbox (payload) VALUES ('{\"n\":\"A\"}'::jsonb)") # id=1
b.execute("INSERT INTO outbox (payload) VALUES ('{\"n\":\"B\"}'::jsonb)") # id=2
b.commit()
a.commit()
with psycopg.connect(DSN, autocommit=True) as r:
changes = r.execute(
"SELECT data FROM pg_logical_slot_get_changes('outbox_bench_slot', NULL, NULL)"
).fetchall()
for (line,) in changes:
print(" ", line)
r.execute("SELECT pg_drop_replication_slot('outbox_bench_slot')")
a.close()
b.close()
print(
"\n Строка с id=2 идёт первой, потому что её транзакция закоммитилась\n"
" первой. Номер транзакции у неё при этом БОЛЬШЕ — значит, порядок\n"
" не по id и не по xid, а по моменту коммита. Тот же порядок, что в\n"
" разделе 4: разница между опросом таблицы и чтением WAL — в цене и\n"
" задержке, а не в том, что доедет."
)
# ---------------------------------------------------------------------------
# 6. Два релея сразу
# ---------------------------------------------------------------------------
def section_skip_locked(c: psycopg.Connection) -> None:
head(6, "Два релея на одной таблице: SKIP LOCKED и что без него")
fresh(c)
cur = c.cursor()
cur.executemany(
"INSERT INTO outbox (payload) VALUES (%s::jsonb)", [(f'{{"n":{i}}}',) for i in range(6)]
)
c.commit()
r1 = psycopg.connect(DSN)
r2 = psycopg.connect(DSN)
q = (
"SELECT id FROM outbox WHERE published_at IS NULL ORDER BY id LIMIT 3 "
"FOR UPDATE SKIP LOCKED"
)
g1 = [i for (i,) in r1.execute(q).fetchall()]
g2 = [i for (i,) in r2.execute(q).fetchall()]
print(f" с SKIP LOCKED: релей 1 взял {g1}, релей 2 взял {g2}")
print(f" пересечение: {set(g1) & set(g2) or 'пусто'}")
r1.rollback()
r2.rollback()
q2 = "SELECT id FROM outbox WHERE published_at IS NULL ORDER BY id LIMIT 3 FOR UPDATE"
r1.execute(q2) # держит блокировку
t0 = time.perf_counter()
try:
r2.execute("SET lock_timeout = '300ms'")
r2.execute(q2)
print(" без SKIP LOCKED: релей 2 прошёл сразу")
except psycopg.errors.LockNotAvailable:
ms = (time.perf_counter() - t0) * 1000
print(f" без SKIP LOCKED: релей 2 ждал и упал по lock_timeout ({ms:.0f} мс)")
print(" то есть второй релей не ускоряет, а простаивает")
r1.rollback()
r2.rollback()
r1.close()
r2.close()
# ---------------------------------------------------------------------------
# 7. Ровно один раз не бывает
# ---------------------------------------------------------------------------
def section_at_least_once(c: psycopg.Connection) -> None:
head(7, "Между «отправил» и «пометил» тоже можно упасть")
fresh(c)
c.execute("INSERT INTO outbox (payload) VALUES ('{\"n\":\"A\"}'::jsonb)")
c.commit()
broker: list[int] = []
with psycopg.connect(DSN) as r:
rows = r.execute(
"SELECT id FROM outbox WHERE published_at IS NULL ORDER BY id FOR UPDATE SKIP LOCKED"
).fetchall()
broker += [i for (i,) in rows]
r.rollback() # падение до UPDATE published_at
print(f" релей отправил {broker}, упал до пометки")
with psycopg.connect(DSN, autocommit=True) as r:
rows = r.execute(
"SELECT id FROM outbox WHERE published_at IS NULL ORDER BY id"
).fetchall()
broker += [i for (i,) in rows]
r.execute("UPDATE outbox SET published_at = now() WHERE id = ANY(%s)", ([i for (i,) in rows],))
print(f" после перезапуска отправил {[i for (i,) in rows]}")
print(f" брокер получил: {broker} — событие ушло дважды")
print(
"\n Поменять местами пометку и отправку нельзя: тогда падение между\n"
" ними потеряет событие, а это хуже. Outbox даёт «хотя бы один раз»\n"
" и не может дать больше — идемпотентность нужна на стороне\n"
" потребителя, и её обычно строят на id события."
)
# ---------------------------------------------------------------------------
# 8. Долгая транзакция задерживает всё, что закоммитилось позже
# ---------------------------------------------------------------------------
def section_long_transaction(c: psycopg.Connection) -> None:
head(8, "Долгая транзакция чужих событий НЕ задерживает — отсюда и разрыв")
fresh(c)
c.commit()
slow = psycopg.connect(DSN)
slow.execute("INSERT INTO outbox (payload) VALUES ('{\"n\":\"slow\"}'::jsonb)")
fast = psycopg.connect(DSN)
fast.execute("INSERT INTO outbox (payload) VALUES ('{\"n\":\"fast\"}'::jsonb)")
fast.commit()
with psycopg.connect(DSN, autocommit=True) as r:
visible = r.execute("SELECT payload->>'n' FROM outbox ORDER BY id").fetchall()
print(f" пока долгая транзакция открыта, релей видит: {[n for (n,) in visible]}")
slow.commit()
with psycopg.connect(DSN, autocommit=True) as r:
visible = r.execute("SELECT payload->>'n' FROM outbox ORDER BY id").fetchall()
print(f" после её коммита: {[n for (n,) in visible]}")
print(
"\n Событие «fast» доехало вовремя — незакоммиченная строка не\n"
" загораживает более поздние. Но релей, который сортирует по id и\n"
" помнит отметку, в разделе 3 на этом и сломался: строка «slow»\n"
" появилась ПОЗАДИ уже прочитанной. Отсюда правило: у релея не\n"
" должно быть памяти о том, до какого id он дошёл."
)
slow.close()
fast.close()
def main() -> None:
print(__doc__.split("ЗАПУСК")[0].strip())
try:
conn = psycopg.connect(DSN)
except psycopg.OperationalError as exc:
sys.exit(f"не подключиться: {exc}\nDSN берётся из OUTBOX_DSN, сейчас: {DSN}")
with conn as c:
v = c.execute("SHOW server_version").fetchone()[0]
print(f"\nPostgreSQL {v}, psycopg {psycopg.__version__}")
section_dual_write(c)
section_atomic(c)
section_high_water_mark(c)
section_status_column(c)
section_logical_decoding(c)
section_skip_locked(c)
section_at_least_once(c)
section_long_transaction(c)
print()
if __name__ == "__main__":
main()