Deep Engineering

MEASUREMENT

bench/isolation/explicit_locking.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/isolation-levels

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

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

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

скрипт что показывает
anomalies.py четыре аномалии на трёх уровнях, воспроизведением: неповторяющееся чтение, фантом, потерянное обновление, write skew — и где именно приходит отказ
cost.py пропускная способность, разброс между кругами и доли отказов 40001 и 40P01 порознь
explicit_locking.py держит ли инвариант SELECT ... FOR UPDATE на слабом уровне и где этот приём молча не срабатывает

Запуск (нужен живой сервер; адрес — в DE_BENCH_DSN):

DE_BENCH_DSN="postgres://postgres@localhost:5433/postgres" \
  python3 bench/isolation/anomalies.py
DE_BENCH_DSN="postgres://postgres@localhost:5433/postgres" \
  python3 bench/isolation/cost.py
DE_BENCH_DSN="postgres://postgres@localhost:5433/postgres" \
  python3 bench/isolation/explicit_locking.py

Почему аномалии воспроизводятся, а не описываются

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

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

Ошибки этих замеров, оставленные в истории

Первая редакция anomalies.py утверждала в заголовке, что PostgreSQL молча подменяет запрошенный READ UNCOMMITTED на READ COMMITTED. Запуск это опроверг: transaction_isolation возвращает ровно read uncommitted — имя уровня сохраняется. Подменяется поведение, и проверять надо его: соседняя транзакция меняет строку и не фиксирует, а мы пробуем это увидеть.

Замечание не о PostgreSQL, а о методе: утверждение про имя настройки и утверждение про поведение — разные утверждения, и первое ничего не говорит о втором. Скрипт теперь печатает оба.

Вторая: cost.py складывал 40001 и 40P01 в один счётчик, а колонку подписывал «отказов 40001». Это разные SQLSTATE — serialization_failure и deadlock_detected, — и подпись величины обязана соответствовать тому, что посчитано, независимо от того, повезло ли в конкретном прогоне. Счётчики разделены; дедлоков за все прогоны не случилось ни одного, и теперь это видно в отчёте явной строкой, а не предполагается.

Третья: пять кругов оказалось мало для того вывода, который на них строился. По лучшему из пяти кругов выходило, что на редких столкновениях SERIALIZABLE опережает REPEATABLE READ. Девять кругов показали, что порядок между этими двумя от круга к кругу переворачивается, а разница между ними меньше собственного разброса каждого. Вывод «SERIALIZABLE впереди» снят как неподтверждённый, а в отчёт добавлены худший круг и строка «впереди в N из 9»: без них число выглядит точнее, чем оно есть.

Протокол замера цены

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

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

Наборы строк сравнивать между собой нельзя. У набора из 4 строк и набора из 64 разная вероятность столкнуться, и числа между ними несопоставимы — сравнивать надо уровни внутри одного набора. Ради этого набора и два: в этой нагрузке плотность соперничества двигает цену сильнее, чем переход между двумя верхними уровнями. Это не значит, что уровень ни при чём: у READ COMMITTED доля отказов нулевая на обоих наборах, у верхних — нет.

Script

264 lines
"""Можно ли удержать инвариант БЕЗ SERIALIZABLE — проверкой, а не рассуждением.

ЗАЧЕМ ЭТОТ СКРИПТ. Соседний anomalies.py показывает, что write skew проходит
на REPEATABLE READ, а SERIALIZABLE его останавливает. Отсюда легко сделать
вывод шире, чем замер: «инвариант охраняет только SERIALIZABLE». Это неверно.
Тот же инвариант можно удержать и на слабом уровне — явной блокировкой. Но
тогда корректность обеспечивает уже не уровень, а выбранная блокировка, и у
неё свои границы. Скрипт показывает и то, и другое: где приём работает и где
он молча не срабатывает.

ЧТО ИМЕННО ПРОВЕРЯЕТСЯ. Два сценария с одинаковой формой «прочитать предикат,
решить, записать» и одинаковой защитой `SELECT ... FOR UPDATE`:

  1. ПРЕДИКАТ ПО СУЩЕСТВУЮЩИМ СТРОКАМ. Двое дежурных, инвариант «останется
     хотя бы один». `FOR UPDATE` захватывает те самые строки, по которым
     принимается решение, — второй транзакции придётся ждать.

  2. ПРЕДИКАТ ПО ОТСУТСТВУЮЩИМ СТРОКАМ. Инвариант «на слот не больше одной
     записи», и записи ещё нет. `FOR UPDATE` блокирует ТОЛЬКО фактически
     возвращённые строки; пустая выборка не блокирует ничего, и обе
     транзакции вставляют свою.

Разница между сценариями и есть граница совета «возьмите FOR UPDATE».

ПОЧЕМУ ЗДЕСЬ НЕТ ПАУЗ-УГАДАЕК. Вторая транзакция обязана реально упереться в
блокировку первой. Вместо sleep скрипт ждёт появления неудовлетворённой
блокировки в `pg_locks` с отдельного соединения — то есть ждёт наступления
самого события, а не отмеренного времени.

ЗАПУСК:

    DE_BENCH_DSN="postgres://postgres@localhost:5433/postgres" \\
      python3 bench/isolation/explicit_locking.py
"""

import os
import sys
import threading
import time

import psycopg

DSN = os.environ.get("DE_BENCH_DSN", "postgres://postgres@localhost:5433/postgres")

LEVELS = ["READ COMMITTED", "REPEATABLE READ", "SERIALIZABLE"]

# В psycopg 3 классы ошибок разложены по SQLSTATE плоско: общего предка у
# 40001 и 40P01 нет, ловить приходится перечислением.
ROLLBACK = (psycopg.errors.SerializationFailure, psycopg.errors.DeadlockDetected)


def run(conn, sql: str) -> None:
    with conn.cursor() as cur:
        cur.execute(sql)


def rows(conn, sql: str) -> list:
    with conn.cursor() as cur:
        cur.execute(sql)
        return cur.fetchall()


def scalar(conn, sql: str):
    got = rows(conn, sql)
    return got[0][0] if got else None


def begin(conn, level: str) -> None:
    run(conn, f"BEGIN ISOLATION LEVEL {level}")


def await_blocked(deadline: float = 5.0) -> bool:
    """Ждать, пока в базе появится ЖДУЩАЯ блокировка.

    Это и есть барьер вместо sleep: событие, которого мы ждём, наблюдаемо
    напрямую. Отдельное соединение, чтобы не участвовать в самой сцене.
    """
    stop = time.perf_counter() + deadline
    with psycopg.connect(DSN, autocommit=True) as watch:
        while time.perf_counter() < stop:
            if scalar(watch, "SELECT count(*) FROM pg_locks WHERE NOT granted") > 0:
                return True
            time.sleep(0.01)
    return False


# ------------------------------------------------------- сценарий 1: дежурные


def setup_oncall() -> None:
    with psycopg.connect(DSN, autocommit=True) as conn:
        run(conn, "DROP TABLE IF EXISTS oncall")
        run(conn, "CREATE TABLE oncall (doctor text primary key, on_call bool not null)")
        run(conn, "INSERT INTO oncall VALUES ('Алиса', true), ('Борис', true)")


def oncall_with_locking(level: str) -> str:
    """Тот же write skew, но обе транзакции берут строки предиката FOR UPDATE.

    FOR UPDATE здесь захватывает ИМЕННО ТЕ строки, по которым принимается
    решение, поэтому вторая транзакция не может прочитать предикат, пока
    первая не закончила.
    """
    setup_oncall()
    outcome: list[str] = []

    def second(conn) -> None:
        try:
            begin(conn, level)
            # Упрётся в блокировку первой транзакции и будет ждать.
            on_call = rows(conn, "SELECT doctor FROM oncall WHERE on_call FOR UPDATE")
            if len(on_call) < 2:
                # Увидела правду: уходить нельзя.
                conn.rollback()
                outcome.append("вторая отказалась уходить")
                return
            run(conn, "UPDATE oncall SET on_call = false WHERE doctor = 'Борис'")
            conn.commit()
            outcome.append("вторая ушла тоже")
        except ROLLBACK as exc:
            conn.rollback()
            outcome.append(f"вторая остановлена, {exc.sqlstate}")

    with psycopg.connect(DSN) as t1, psycopg.connect(DSN) as t2:
        t1.autocommit = t2.autocommit = False

        begin(t1, level)
        first = rows(t1, "SELECT doctor FROM oncall WHERE on_call FOR UPDATE")
        if len(first) != 2:
            return f"негодная подготовка: {len(first)}"

        thread = threading.Thread(target=second, args=(t2,))
        thread.start()
        if not await_blocked():
            # Если вторая НЕ заблокировалась — приём не сработал, и это тоже
            # результат, который надо показать, а не спрятать.
            outcome.append("вторая не заблокировалась")

        run(t1, "UPDATE oncall SET on_call = false WHERE doctor = 'Алиса'")
        t1.commit()
        thread.join(timeout=10)

    with psycopg.connect(DSN, autocommit=True) as check:
        left = scalar(check, "SELECT count(*) FROM oncall WHERE on_call")
    note = "; ".join(outcome) if outcome else "нет исхода"
    return f"дежурных {left}{note}"


# ------------------------------------------------- сценарий 2: пустая выборка


def setup_slots() -> None:
    with psycopg.connect(DSN, autocommit=True) as conn:
        run(conn, "DROP TABLE IF EXISTS bookings")
        run(conn, "CREATE TABLE bookings (id serial primary key, slot int not null, who text not null)")


def empty_predicate_with_locking(level: str) -> str:
    """Инвариант «на слот не больше одной записи», записи ещё нет.

    Обе транзакции защищаются тем же приёмом: смотрят занятость слота
    `FOR UPDATE` и вставляют, если пусто. Но блокировать нечего: строк,
    удовлетворяющих условию, не существует.
    """
    setup_slots()
    outcome: list[str] = []

    def second(conn) -> None:
        try:
            begin(conn, level)
            taken = rows(conn, "SELECT id FROM bookings WHERE slot = 1 FOR UPDATE")
            if taken:
                conn.rollback()
                outcome.append("вторая увидела занятость")
                return
            run(conn, "INSERT INTO bookings (slot, who) VALUES (1, 'вторая')")
            conn.commit()
            outcome.append("вторая вставила")
        except ROLLBACK as exc:
            conn.rollback()
            outcome.append(f"вторая остановлена, {exc.sqlstate}")

    with psycopg.connect(DSN) as t1, psycopg.connect(DSN) as t2:
        t1.autocommit = t2.autocommit = False

        begin(t1, level)
        rows(t1, "SELECT id FROM bookings WHERE slot = 1 FOR UPDATE")  # пусто

        # Обе читают предикат ДО того, как хоть одна вставила: блокировки,
        # которая заставила бы вторую ждать, здесь не возникает.
        thread = threading.Thread(target=second, args=(t2,))
        thread.start()
        time.sleep(0.05)  # дать второй дойти до своего чтения; она не блокируется

        # Отказ может прийти и на INSERT, а не только на COMMIT: под SSI вторая
        # транзакция успевает зафиксироваться и создать зависимость раньше.
        # Поэтому в защите обе команды, а не одна фиксация, — ровно то, о чём
        # говорит правило «повторять транзакцию целиком».
        try:
            run(t1, "INSERT INTO bookings (slot, who) VALUES (1, 'первая')")
            t1.commit()
        except ROLLBACK as exc:
            t1.rollback()
            outcome.append(f"первая остановлена на своей команде, {exc.sqlstate}")
        thread.join(timeout=10)

    with psycopg.connect(DSN, autocommit=True) as check:
        n = scalar(check, "SELECT count(*) FROM bookings WHERE slot = 1")
    note = "; ".join(outcome) if outcome else "нет исхода"
    return f"записей на слот {n}{note}"


def main() -> None:
    with psycopg.connect(DSN, autocommit=True) as conn:
        version = scalar(conn, "SHOW server_version")

    print(f"PostgreSQL {version} | инвариант без SERIALIZABLE: где явная блокировка держит")
    print()

    checks = [
        ("предикат по существующим строкам", oncall_with_locking),
        ("предикат по отсутствующим строкам", empty_predicate_with_locking),
    ]
    header = "сценарий, везде SELECT ... FOR UPDATE"
    width = max([len(name) for name, _ in checks] + [len(header)])

    print(f"  {header:<{width}}  " + "  ".join(f"{l:<46}" for l in LEVELS))
    for name, fn in checks:
        cells = []
        for level in LEVELS:
            try:
                cells.append(fn(level))
            except Exception as exc:  # noqa: BLE001 — печатаем любой отказ как есть
                cells.append(type(exc).__name__)
        print(f"  {name:<{width}}  " + "  ".join(f"{c:<46}" for c in cells))

    print()
    print("ЧТО ИЗ ЭТОГО СЛЕДУЕТ")
    print("  «Инвариант охраняет только SERIALIZABLE» — слишком широко. В первом")
    print("  сценарии явная блокировка удерживает его и на слабом уровне: строки")
    print("  предиката существуют, FOR UPDATE их захватывает, вторая транзакция")
    print("  ждёт и видит правду.")
    print()
    print("  Но охраняет тогда не уровень, а выбранная блокировка — и вместе с")
    print("  уровнем она даёт разный исход: на READ COMMITTED вторая транзакция")
    print("  дочитывает свежую строку и отступает сама, на верхних уровнях тот же")
    print("  приём оборачивается отказом 40001. Оба исхода инвариант сохраняют,")
    print("  но обрабатывать их приложению приходится по-разному.")
    print()
    print("  Во втором сценарии тот же приём не срабатывает вовсе: FOR UPDATE")
    print("  блокирует только фактически возвращённые строки, а их нет. Пустая")
    print("  выборка не создаёт блокировку диапазона. Здесь помогает либо")
    print("  SERIALIZABLE, либо ограничение уникальности — то есть защита,")
    print("  которая не зависит от того, что транзакция успела прочитать.")


if __name__ == "__main__":
    try:
        main()
    except psycopg.OperationalError as exc:
        print(f"нет соединения с базой: {exc}", file=sys.stderr)
        print(f"адрес: {DSN}", file=sys.stderr)
        raise SystemExit(1)