Event store на Postgres: sqlx, транзакция, оптимистичная блокировка
открытый урокЭтот раздел читается без входа. Войди, чтобы отмечать прогресс, вести заметки и решать задачи в редакторе. войти
Event store на Postgres: sqlx, транзакция, оптимистичная блокировка
Сцена · события первого класса
Decider у нас чистый и без памяти: дай ему историю, он скажет, что делать дальше. Теперь дадим ему память. В event sourcing хранилище, это не таблица с текущим снимком брони, которую мы перезаписываем. Это append-only журнал событий: новые дописываются в конец, старые не меняются никогда. Один стрим, это одна бронь. Сегодня дадим Decider два хранилища с одним контрактом: in-memory для тестов и Postgres для прода.
Карта урока
- Контракт
EventStore. Один асинхронный трейт, две реализации. - Оптимистичная блокировка.
expected_versionи почему две команды на одну бронь не затрут друг друга. - In-memory store. Простейшая реализация на
Mutex, для тестов и проекций. - Схема Postgres. Таблица
events,jsonb, иunique (stream_id, version). - sqlx без живой базы. Почему runtime query-API, а не макрос.
- 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 свежим без опроса базы.