← назад к разделу

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

Обязательно

Где граница транзакции

Граница транзакции совпадает с границей единицы работы: один обработчик команды, одна транзакция. Самый ясный способ записать это:

def pay_order(order_id: int, amount: Decimal) -> None:
    with SessionFactory() as session, session.begin():
        order = session.get(Order, order_id)
        order.pay(amount)

session.begin() открывает транзакцию, на выходе из блока делает commit, при исключении rollback. Без begin() сессия всё равно откроет транзакцию при первом запросе (autobegin), но тогда commit придётся звать руками, и забытый commit молча откатится при закрытии сессии: данные «не сохранились», ошибки нет. Явный begin() снимает этот класс проблем.

Два правила, которые из этого следуют. Внутри транзакции нет сетевых вызовов: HTTP к соседнему сервису, отправка в Kafka, запись в S3 делаются до или после, иначе соединение и блокировки держатся на время чужой задержки. И транзакция не пересекает границу обработчика: репозиторий не делает commit сам, он только читает и пишет в сессию, которую ему дали; иначе половина команды зафиксируется, а вторая откатится.

Во FastAPI сессию выдаёт зависимость, и транзакцию обычно открывают там же, чтобы обработчик не думал о ней, как показано в статье про SQLAlchemy во FastAPI. В обработчике CQRS границу держит сам обработчик команды, об этом статья про command side.

Savepoint: откатить часть, не всё

Иногда внутри транзакции есть шаг, который может не получиться, и это не повод откатывать всё: вставка записи, которая, возможно, уже есть. begin_nested() открывает savepoint:

with session.begin():
    order = session.get(Order, order_id)
    order.confirm()
    try:
        with session.begin_nested():
            session.add(AuditEntry(order_id=order_id, action="confirm"))
            session.flush()
    except IntegrityError:
        pass
    session.add(order)

При исключении внутри вложенного блока откатывается только savepoint, а внешняя транзакция живёт дальше. Без savepoint IntegrityError переводит всё соединение PostgreSQL в состояние «транзакция прервана», и любой следующий запрос падает до rollback. Поэтому «поймать IntegrityError и продолжить» без begin_nested не работает в принципе.

Для «вставить, если нет» у PostgreSQL есть более прямой путь: insert(...).on_conflict_do_nothing() из sqlalchemy.dialects.postgresql, без исключений и savepoint.

Уровень изоляции

PostgreSQL по умолчанию работает в READ COMMITTED: каждый запрос видит данные, зафиксированные к его началу. Для большинства команд этого достаточно, если изменения защищены блокировкой строки или версией (о них ниже). Когда нужен снимок на всю транзакцию (отчёт из нескольких запросов, которые должны сойтись), уровень поднимают до REPEATABLE READ на одно соединение:

with engine.connect().execution_options(isolation_level="REPEATABLE READ") as conn:
    with Session(bind=conn) as session, session.begin():
        ...

Или на весь движок через create_engine(..., isolation_level="REPEATABLE READ"), но тогда каждый обработчик должен быть готов к ошибке сериализации при одновременном изменении и уметь повторить попытку. Что именно даёт каждый уровень и какие аномалии остаются, подробно разбирает статья про уровни изоляции PostgreSQL; здесь важно одно: уровень изоляции не заменяет блокировку строки для сценария «прочитал, посчитал, записал».

Потерянное обновление и два способа от него защититься

Два обработчика одновременно читают остаток на счёте 100, каждый вычитает 30 и записывает 70. Одно списание потеряно. В READ COMMITTED это штатное поведение, и ORM от него не спасает: order.balance -= 30 это чтение в Python и UPDATE ... SET balance = 70 без оглядки на других.

Пессимистичная блокировка: взять строку под FOR UPDATE при чтении, тогда второй обработчик подождёт первого и прочитает уже 70:

order = session.scalars(
    select(Order).where(Order.id == order_id).with_for_update()
).one()
order.withdraw(30)

with_for_update(nowait=True) вместо ожидания сразу падает, если строка занята; with_for_update(skip_locked=True) пропускает занятые строки, и это основа очередей на таблице: воркеры берут «следующую свободную задачу» без дублей. С joinedload в одном запросе PostgreSQL откажется: FOR UPDATE нельзя применить к внешней стороне LEFT OUTER JOIN, поэтому связи догружают отдельно или блокируют только нужную таблицу через with_for_update(of=Order).

Оптимистичная версия: колонка версии в таблице, которую SQLAlchemy проверяет при каждом UPDATE:

class Order(Base):
    __tablename__ = "orders"
    id: Mapped[int] = mapped_column(primary_key=True)
    balance: Mapped[Decimal]
    version: Mapped[int] = mapped_column(default=1)

    __mapper_args__ = {"version_id_col": version}

UPDATE orders SET balance = ?, version = 2 WHERE id = ? AND version = 1. Если кто-то успел раньше, строка не совпала, и flush поднимает StaleDataError. Обработчик превращает её в ответ 409 и просит клиента повторить, как описано в статье про ошибки API. Проверено на SQLAlchemy 2.1: две сессии загрузили один заказ, первая зафиксировала изменение, вторая упала на commit с StaleDataError.

Выбор между ними: блокировка строки для коротких транзакций с высокой конкуренцией за одну строку (остаток, счётчик), версия для агрегатов, которые редактируют люди, где конфликт редок и повтор дёшев. Для агрегата из нескольких таблиц версия на корне защищает весь агрегат, если любое изменение проходит через корень и поднимает его версию.

Что транзакция не защищает

Транзакция не делает атомарной пару «записать в базу и отправить событие». Если commit прошёл, а отправка в Kafka упала, событие потеряно; если отправили до commit, а он откатился, событие ложное. Ответ на это outbox: событие пишут в таблицу той же транзакцией, а отправляет его отдельный процесс. Подробно об этом статья про распределённые паттерны.

И транзакция не защищает от длинного with session.begin(), внутри которого ждут пользователя, таймер или внешний сервис. Такие блоки держат соединение из пула и строки под блокировкой, и под нагрузкой пул заканчивается первым.

Дополнительно: при первом чтении можно пропустить

Глубже: вложенные вызовы и одна транзакция на всехрасширенное

В Java-мире привычна аннотация, которая «присоединяется к текущей транзакции или открывает новую». В SQLAlchemy такой магии нет, и это к лучшему: транзакцию держит тот, кто создал сессию, а все вызываемые функции получают сессию аргументом и в ней работают. Повторный session.begin() внутри уже открытой транзакции поднимает InvalidRequestError, и это сигнал, что функция пытается управлять границей, которая ей не принадлежит.

Если функция должна работать и сама по себе, и внутри чужой транзакции, её делят на две: внутренняя принимает сессию и ничего не фиксирует, внешняя создаёт сессию, открывает begin() и зовёт внутреннюю. Так репозитории и доменные сервисы всегда «внутренние», а граница одна, в обработчике команды.

Отдельный случай: join_transaction_mode у сессии, созданной поверх внешнего Connection. В тестах соединение открывают с транзакцией, сессию привязывают к нему, и после теста откатывают всё разом, включая commit, которые сделал код: режим create_savepoint превращает commit сессии в освобождение savepoint. Это самый быстрый способ изолировать тесты на реальной базе, и его показывает статья про интеграционные тесты.

Коротко

  • Граница транзакции равна границе обработчика команды: with SessionFactory() as session, session.begin(); репозиторий сам commit не делает.
  • Внутри транзакции нет сетевых вызовов и ожиданий; соединение и блокировки держатся до конца блока.
  • begin_nested() даёт savepoint: после IntegrityError без него соединение PostgreSQL прервано до rollback; для «вставить, если нет» лучше on_conflict_do_nothing.
  • READ COMMITTED по умолчанию; REPEATABLE READ для снимка на несколько запросов через execution_options(isolation_level=...), с готовностью к повтору.
  • Потерянное обновление лечат with_for_update() (короткие транзакции, одна горячая строка) или version_id_col со StaleDataError в 409 (редкие конфликты людей).
  • skip_locked=True это очередь на таблице без дублей между воркерами.
  • Запись в базу и отправка события не атомарны без outbox.
  • Повторный begin() внутри транзакции это ошибка; функции делят на внутреннюю с сессией и внешнюю с границей.

Что почитать дальше