Раздел 23 · Rust

Event store на Postgres: sqlx, транзакция, оптимистичная блокировка

senior~50 мин

открытый урокЭтот раздел читается без входа. Войди, чтобы отмечать прогресс, вести заметки и решать задачи в редакторе. войти

Event store на Postgres: sqlx, транзакция, оптимистичная блокировка

Сцена · события первого класса

Decider у нас чистый и без памяти: дай ему историю, он скажет, что делать дальше. Теперь дадим ему память. В event sourcing хранилище, это не таблица с текущим снимком брони, которую мы перезаписываем. Это append-only журнал событий: новые дописываются в конец, старые не меняются никогда. Один стрим, это одна бронь. Сегодня дадим Decider два хранилища с одним контрактом: in-memory для тестов и Postgres для прода.

Карта урока

  1. Контракт EventStore. Один асинхронный трейт, две реализации.
  2. Оптимистичная блокировка. expected_version и почему две команды на одну бронь не затрут друг друга.
  3. In-memory store. Простейшая реализация на Mutex, для тестов и проекций.
  4. Схема Postgres. Таблица events, jsonb, и unique (stream_id, version).
  5. sqlx без живой базы. Почему runtime query-API, а не макрос.
  6. append в транзакции. INSERT, перехват unique violation, перевод в доменную ошибку.

Раздел 1 · Контракт EventStore

Хранилищу нужны ровно две операции: прочитать стрим и дописать в него события. Опишем их трейтом. Реальная реализация ходит в Postgres, то есть асинхронна, поэтому трейт async через async-trait, и in-memory тоже async, чтобы интерфейс был один:

use async_trait::async_trait;

/// Срез стрима: его длина (версия) и события по порядку.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct StreamSlice {
    pub version: u32,
    pub events: Vec<ReservationEvent>,
}

#[async_trait]
pub trait EventStore {
    /// Прочитать стрим целиком. Пустой стрим, это version=0 и нет событий.
    async fn load(&self, stream: &ReservationId) -> StreamSlice;

    /// Дописать события при условии, что версия стрима равна expected_version.
    /// Иначе DomainError::ConcurrencyConflict.
    async fn append(
        &self,
        stream: &ReservationId,
        expected_version: u32,
        events: Vec<ReservationEvent>,
    ) -> Result<(), DomainError>;
}

Связка с Decider простая и циклична: load отдаёт StreamSlice, мы делаем replay его событий в состояние, зовём decide, и новые события отправляем в append с той версией, что видели при чтении. Вот это «с той версией, что видели», и есть защита от гонки.

Раздел 2 · Оптимистичная блокировка

Представь две команды на одну бронь одновременно: «засели гостя» и «отмени». Обе прочитали стрим версии 1, обе решили на состоянии Reserved, обе пишут событие версии 2. Без защиты одна затрёт другую, и бронь окажется в противоречивом состоянии.

Решение, это оптимистичная блокировка: append принимает expected_version, и если текущая длина стрима ей не равна, запись отвергается. В терминах нашего трейта это DomainError::ConcurrencyConflict. Тогда первая команда запишется (версия совпала), а вторая получит конфликт, перечитает свежий стрим и перерешает на актуальном состоянии, где, возможно, заселять уже некого. Блокировку мы не держим, только проверяем версию в момент записи, отсюда «оптимистичная».

Раздел 3 · In-memory store

Прежде чем Postgres, реализуем контракт на HashMap под Mutex. Этого хватает для всех BDD-сценариев и для проекций следующего урока, и Postgres ради проверки логики таскать не нужно:

pub struct InMemoryEventStore {
    streams: Mutex<HashMap<ReservationId, Vec<ReservationEvent>>>,
    tx: broadcast::Sender<AllEventsItem>,
}

Поле tx это broadcast-канал: каждое записанное событие транслируется подписчикам, на этом в уроке 4 построится live-проекция. Сам append проверяет версию и дописывает:

async fn append(
    &self,
    stream: &ReservationId,
    expected_version: u32,
    events: Vec<ReservationEvent>,
) -> Result<(), DomainError> {
    {
        let mut streams = self.streams.lock().await;
        let existing = streams.entry(stream.clone()).or_default();
        let actual = existing.len() as u32;
        if actual != expected_version {
            return Err(DomainError::ConcurrencyConflict {
                reservation_id: stream.as_str().to_owned(),
                expected: expected_version,
                actual,
            });
        }
        existing.extend(events.iter().cloned());
    }
    // Транслируем уже после успешной записи и вне блокировки Mutex.
    for event in events {
        let _ = self.tx.send(AllEventsItem { stream_id: stream.clone(), event });
    }
    Ok(())
}

Та же оптимистичная блокировка, что будет в Postgres, только на длине вектора. Заметь: рассылку делаем после записи и за пределами блокировки, чтобы не держать Mutex дольше нужного.

Раздел 4 · Схема Postgres

В проде стрим живёт в таблице. Минимальная схема: id для глобального порядка, stream_id и version для самого стрима, тип и payload для события:

pub const SCHEMA_DDL: &str = "\
create table if not exists events (
  id            bigserial primary key,
  stream_id     text        not null,
  version       integer     not null,
  type          text        not null,
  payload       jsonb       not null,
  created_at    timestamptz not null default now(),
  unique (stream_id, version)
);
create index if not exists events_stream_id_idx on events (stream_id, version);
create index if not exists events_id_idx on events (id);
";

Два решения тут несущие. Во-первых, payload это jsonb: событие сериализуется тем самым внутренне-тегированным serde, и {"type": ...} ложится в колонку как есть, round-trip в обе стороны одной парой функций. Во-вторых, unique (stream_id, version) бесплатно даёт оптимистичную блокировку: вставить дважды одну версию стрима база не позволит.

Раздел 5 · sqlx без живой базы

Подключаемся через sqlx, но есть важная развилка. У sqlx есть макрос query!, который проверяет SQL на этапе компиляции, дёргая настоящую базу при сборке. Для учебного пакета это плохо: cargo build и cargo test не должны требовать поднятый Postgres. Поэтому берём runtime query-API sqlx::query(...), где биндинги и типы пишем руками, зато база на сборке не нужна:

async fn load(&self, stream: &ReservationId) -> StreamSlice {
    let rows = sqlx::query(
        "select version, payload from events where stream_id = $1 order by version asc",
    )
    .bind(stream.as_str())
    .fetch_all(&self.pool)
    .await
    .expect("load events failed");

    let mut events = Vec::with_capacity(rows.len());
    let mut version: u32 = 0;
    for row in rows {
        let payload: serde_json::Value = row.get("payload");
        let raw_version: i32 = row.get("version");
        let event: ReservationEvent =
            serde_json::from_value(payload).expect("payload is not a ReservationEvent");
        events.push(event);
        version = raw_version as u32;
    }
    StreamSlice { version, events }
}

load читает строки по stream_id в порядке версии и десериализует jsonb обратно в ReservationEvent. Тип события восстанавливается по полю type внутри payload, ровно так это и было задумано в прошлом уроке.

Раздел 6 · append в транзакции

Запись нескольких событий должна быть атомарной: либо все, либо ни одного. Оборачиваем INSERT-ы в транзакцию и переводим unique violation в доменный конфликт:

async fn append(
    &self,
    stream: &ReservationId,
    expected_version: u32,
    events: Vec<ReservationEvent>,
) -> Result<(), DomainError> {
    if events.is_empty() {
        return Ok(());
    }
    let mut tx = self.pool.begin().await.expect("begin tx failed");
    for (index, event) in events.iter().enumerate() {
        let version = expected_version + index as u32 + 1;
        let payload = serde_json::to_value(event).expect("event is not serializable");
        let result = sqlx::query(
            "insert into events (stream_id, version, type, payload) values ($1, $2, $3, $4)",
        )
        .bind(stream.as_str())
        .bind(version as i32)
        .bind(event.type_name())
        .bind(payload)
        .execute(&mut *tx)
        .await;

        if let Err(error) = result {
            let _ = tx.rollback().await;
            if is_unique_violation(&error) {
                return Err(DomainError::ConcurrencyConflict {
                    reservation_id: stream.as_str().to_owned(),
                    expected: expected_version,
                    actual: expected_version + events.len() as u32,
                });
            }
            panic!("insert event failed: {error}");
        }
    }
    tx.commit().await.expect("commit tx failed");
    Ok(())
}

Версия каждого события считается от expected_version, поэтому параллельный писатель, успевший раньше, займёт нашу версию, и наш INSERT нарушит unique (stream_id, version). Postgres вернёт ошибку с кодом 23505, и мы распознаём её по SQLSTATE:

/// Это unique violation на паре (stream_id, version)?
fn is_unique_violation(error: &sqlx::Error) -> bool {
    matches!(
        error.as_database_error().and_then(|db| db.code()),
        Some(code) if code == "23505"
    )
}

Конфликт версий, который в in-memory store был сравнением длин, в Postgres прилетает от индекса уникальности, и обе реализации возвращают один и тот же DomainError::ConcurrencyConflict. Контракт EventStore один, поведение совпадает, а вызывающий код не знает и не должен знать, какое хранилище под ним.

Чек-лист

Чек-листготово

Домашка

Дальше

Финальный урок трека, Async-проекции и CQRS. Журнал событий есть, и из него можно построить любую модель для чтения: список заселённых, занятость по этажам, выручку за день. Соберём live-проекцию на tokio::spawn, broadcast и watch, которая держит read-model свежим без опроса базы.