Deep Engineering

MEASUREMENT

bench/outbox/cost.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

271 lines
"""Transaction outbox: сколько это стоит на настоящем PostgreSQL.

Четыре блока, и каждый отвечает на вопрос, который задают, решая, брать
outbox или нет:

  1. Сколько стоит лишняя строка в бизнес-транзакции — и сколько стоит
     обойтись без неё, записав событие отдельной транзакцией.
  2. Сколько строк в секунду вытягивает релей и от чего это зависит.
  3. Сколько стоит холостой опрос — то, чем релей занят почти всё время.
  4. Сколько стоит уборка и сколько места остаётся после неё.

ЧТО ЗДЕСЬ СРАВНИМО. Блок 1 — три способа сделать ОДНУ И ТУ ЖЕ работу,
поэтому их числа сравниваются между собой; они и измеряются в одном прогоне,
чередуясь, чтобы состояние базы было общим. Блоки 2–4 сравнивать с блоком 1
нельзя: там другая работа и другие единицы.

ПОЧЕМУ ЧИСЛА БУДУТ ДРУГИМИ. Всё, что здесь измеряется, упирается в fsync при
коммите, а он зависит от диска. Абсолютные значения поэтому не переносятся с
машины на машину; переносятся отношения внутри блока 1 и порядок величин в
блоках 2–4.

ЗАПУСК:
    createdb -p 5433 outbox_bench
    python3 bench/outbox/cost.py

Снято на PostgreSQL 16.13 (synchronous_commit = on, fsync = on,
wal_sync_method = fdatasync), psycopg 3.3.4, Python 3.11.15.
"""

from __future__ import annotations

import os
import statistics
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")
PAYLOAD = '{"type":"OrderPlaced","order_id":1,"total":4200}'


def head(n: int, title: str) -> None:
    print(f"\n{'=' * 72}\n{n}. {title}\n{'=' * 72}")


def make_tables(c: psycopg.Connection, *, partial_index: bool = True) -> None:
    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)"
    )
    if partial_index:
        c.execute("CREATE INDEX outbox_unpub ON outbox (id) WHERE published_at IS NULL")


# ---------------------------------------------------------------------------
# 1. Цена строки outbox в бизнес-транзакции
# ---------------------------------------------------------------------------
def block_transaction_cost(c: psycopg.Connection) -> None:
    head(1, "Цена события: лишняя строка против лишней транзакции")

    make_tables(c)
    c.commit()
    cur = c.cursor()

    def plain() -> None:
        cur.execute("INSERT INTO orders (total) VALUES (4200)")
        c.commit()

    def with_outbox() -> None:
        cur.execute("INSERT INTO orders (total) VALUES (4200)")
        cur.execute("INSERT INTO outbox (payload) VALUES (%s::jsonb)", (PAYLOAD,))
        c.commit()

    def two_transactions() -> None:
        cur.execute("INSERT INTO orders (total) VALUES (4200)")
        c.commit()
        cur.execute("INSERT INTO outbox (payload) VALUES (%s::jsonb)", (PAYLOAD,))
        c.commit()

    cases = {
        "заказ без события": plain,
        "заказ + outbox, одна транзакция": with_outbox,
        "заказ и событие двумя транзакциями": two_transactions,
    }

    for fn in cases.values():
        for _ in range(300):
            fn()

    # Чередование: все три способа в одном прогоне, по кругу. Иначе третий
    # мерился бы на таблице, в которой уже лежит вдвое больше строк.
    samples: dict[str, list[float]] = {k: [] for k in cases}
    for _ in range(1500):
        for name, fn in cases.items():
            t0 = time.perf_counter()
            fn()
            samples[name].append((time.perf_counter() - t0) * 1e6)

    base = statistics.median(samples["заказ без события"])
    print("  медиана и 90-й перцентиль, мкс на один заказ:\n")
    for name, xs in samples.items():
        xs.sort()
        med = statistics.median(xs)
        p90 = xs[int(len(xs) * 0.9)]
        print(f"    {name:<36} {med:8.1f}  p90 {p90:8.1f}   x{med / base:.2f}")

    print(
        "\n  Строка outbox добавляет к транзакции проценты: она пишется в тот\n"
        "  же WAL и коммитится тем же fsync, что и заказ. Отдельная\n"
        "  транзакция для события стоит дороже вдвое — и это второй fsync.\n"
        "  То есть небезопасный способ здесь ещё и медленнее безопасного."
    )


# ---------------------------------------------------------------------------
# 2. Пропускная способность релея
# ---------------------------------------------------------------------------
def block_relay_throughput(c: psycopg.Connection) -> None:
    head(2, "Пропускная способность релея и размер партии")

    rows_total = int(os.environ.get("OUTBOX_ROWS", "20000"))
    cur = c.cursor()
    print(f"  {rows_total} строк, каждый раз таблица наполняется заново\n")

    for batch in (1, 10, 100, 1000, 5000):
        make_tables(c)
        with cur.copy("COPY outbox (payload) FROM STDIN") as cp:
            for _ in range(rows_total):
                cp.write_row([PAYLOAD])
        c.commit()

        t0 = time.perf_counter()
        done = 0
        while True:
            got = cur.execute(
                "SELECT id FROM outbox WHERE published_at IS NULL"
                " ORDER BY id LIMIT %s FOR UPDATE SKIP LOCKED",
                (batch,),
            ).fetchall()
            if not got:
                break
            cur.execute(
                "UPDATE outbox SET published_at = now() WHERE id = ANY(%s)",
                ([i for (i,) in got],),
            )
            c.commit()
            done += len(got)
        dt = time.perf_counter() - t0
        print(
            f"    партия {batch:>5}: {dt:7.2f} с  {done / dt:9.0f} строк/с"
            f"  {dt / done * 1e6:7.1f} мкс на строку"
        )

    print(
        "\n  Отправки в брокер здесь нет вовсе — измеряется только работа с\n"
        "  базой. Разница между крайними строками почти вся состоит из\n"
        "  коммитов: партия по одной строке платит fsync за каждое событие."
    )


# ---------------------------------------------------------------------------
# 3. Холостой опрос
# ---------------------------------------------------------------------------
def block_idle_poll(c: psycopg.Connection) -> None:
    head(3, "Холостой опрос: чем релей занят почти всё время")

    filled = int(os.environ.get("OUTBOX_FILL", "500000"))
    cur = c.cursor()
    query = "SELECT id FROM outbox WHERE published_at IS NULL ORDER BY id LIMIT 100"

    def median_us(times: int) -> float:
        xs = []
        for _ in range(times):
            t0 = time.perf_counter()
            cur.execute(query)
            cur.fetchall()
            xs.append((time.perf_counter() - t0) * 1e6)
        return statistics.median(xs)

    make_tables(c, partial_index=False)
    print(f"  наполняю {filled} УЖЕ ОТПРАВЛЕННЫХ строк...")
    with cur.copy("COPY outbox (payload, published_at) FROM STDIN") as cp:
        for _ in range(filled):
            cp.write_row([PAYLOAD, "2026-08-01 00:00:00+00"])
    c.commit()
    c.execute("ANALYZE outbox")
    c.commit()
    size = c.execute("SELECT pg_size_pretty(pg_total_relation_size('outbox'))").fetchone()[0]
    print(f"  размер таблицы: {size}\n")

    without = median_us(30)
    print(f"    без частичного индекса: {without:10.1f} мкс")

    c.execute("CREATE INDEX outbox_unpub ON outbox (id) WHERE published_at IS NULL")
    c.execute("ANALYZE outbox")
    c.commit()
    with_index = median_us(300)
    idx = c.execute("SELECT pg_size_pretty(pg_relation_size('outbox_unpub'))").fetchone()[0]
    print(f"    с частичным индексом:   {with_index:10.1f} мкс   (x{without / with_index:.0f})")
    print(f"    размер частичного индекса при пустой очереди: {idx}")

    print(
        "\n  Индекс построен по условию «не отправлено», поэтому в нём лежат\n"
        "  только неотправленные строки — когда очередь пуста, он пуст тоже.\n"
        "  Полный индекс по published_at такого не умеет: он хранит запись\n"
        "  для каждой строки таблицы и растёт вместе с ней."
    )


# ---------------------------------------------------------------------------
# 4. Уборка
# ---------------------------------------------------------------------------
def block_cleanup(c: psycopg.Connection) -> None:
    head(4, "Уборка: DELETE, место на диске и VACUUM")

    def size() -> str:
        return c.execute("SELECT pg_size_pretty(pg_total_relation_size('outbox'))").fetchone()[0]

    print(f"  таблица перед уборкой: {size()}")
    t0 = time.perf_counter()
    cur = c.execute("DELETE FROM outbox WHERE published_at < now() - interval '1 hour'")
    dt = time.perf_counter() - t0
    print(f"  DELETE {cur.rowcount} отправленных строк: {dt:.2f} с")
    print(f"  размер сразу после DELETE: {size()}")

    t0 = time.perf_counter()
    c.execute("VACUUM outbox")
    print(f"  VACUUM: {time.perf_counter() - t0:.2f} с, размер: {size()}")

    print(
        "\n  DELETE не возвращает место: строки помечены мёртвыми, файл прежний.\n"
        "  VACUUM возвращает страницы под повторную запись, но отдаёт их\n"
        "  операционной системе только с конца файла — поэтому остаток\n"
        "  ненулевой. Отсюда обычный совет: секции по времени и DROP\n"
        "  вместо DELETE, если через outbox идёт заметный поток."
    )


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:
        c.autocommit = True
        v = c.execute("SHOW server_version").fetchone()[0]
        sc = c.execute("SHOW synchronous_commit").fetchone()[0]
        fs = c.execute("SHOW fsync").fetchone()[0]
        print(f"\nPostgreSQL {v}, synchronous_commit = {sc}, fsync = {fs}")
        c.autocommit = False
        block_transaction_cost(c)
        block_relay_throughput(c)
        c.commit()  # блоки 3–4 работают вне транзакции: там VACUUM
        c.autocommit = True
        block_idle_poll(c)
        block_cleanup(c)
    print()


if __name__ == "__main__":
    main()