Deep Engineering

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'] при id 1 и 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.

Источники

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()