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

В большинстве приложений мы храним текущее состояние: строка в таблице orders со статусом paid и суммой 1000. Каждое изменение — это UPDATE, который перезатирает прошлое: был статус new, стал paid, и старое значение исчезло навсегда. Обычно этого достаточно. Но иногда важно не только что сейчас, а как мы к этому пришли. Event Sourcing отвечает именно на второй вопрос: вместо текущего состояния он хранит неизменяемый поток событий, а состояние вычисляет из него.

Обязательно

Состояние против потока событий

Представьте банковский счёт. Есть два способа рассказать о нём.

Первый — назвать баланс: «на счёте 1200 рублей». Это состояние. Компактно, но вы не знаете, откуда взялась сумма.

Второй — показать выписку: «+1000 зарплата, −300 продукты, +500 возврат». Это поток событий. Баланс здесь нигде не записан отдельно — он получается сложением строк выписки. И именно выписка первична: банк хранит историю операций, а баланс — производная от неё.

Event Sourcing переносит эту идею в код. Мы не храним «заказ оплачен, сумма 1000». Мы храним факты, которые с заказом происходили, в том порядке, в котором они происходили. А текущий заказ собираем, «проигрывая» эти факты один за другим.

Как выглядит поток событий

События — это факты в прошедшем времени: не команда «оплати заказ», а свершившийся факт «заказ оплачен». Для одного заказа поток может выглядеть так:

OrderCreated     { orderId: 42, customer: "Аня" }
ItemAdded        { orderId: 42, lineId: 1, item: "книга", price: 500 }
ItemAdded        { orderId: 42, lineId: 2, item: "ручка", price: 100 }
ItemRemoved      { orderId: 42, lineId: 2 }
OrderPaid        { orderId: 42, amount: 500 }

Обратите внимание на lineId. Событие обязано однозначно указывать, к чему оно применяется, — поэтому ItemRemoved ссылается на номер позиции, а не на название товара. Если бы в заказе было две одинаковые ручки, «убрать ручку» перестало бы быть однозначным: непонятно, какую именно, и повторное проигрывание потока могло бы дать разный результат.

Каждая строка неизменяема: записав событие, мы его больше не трогаем. Если клиент передумал и убрал товар — мы не удаляем ItemAdded, а дописываем новое событие ItemRemoved. История не переписывается, она только растёт. Это ключевое свойство: прошлое нельзя стереть, можно только добавить новый факт.

Хранилище таких событий называется event store. По сути это таблица «только на добавление» (append-only): новые события пишутся в конец, старые никогда не меняются и не удаляются.

Держится это не на честном слове. У каждого события есть номер в своём потоке — (orderId, version), — и на эту пару вешают уникальный индекс. Отсюда два следствия. Во-первых, событие нельзя переписать: попытка записать номер, который уже занят, просто не пройдёт. Во-вторых, так решается вопрос одновременной записи: две команды прочитали заказ на версии 5, обе хотят дописать шестое событие — одна запишет, вторая получит отказ по уникальности, перечитает поток и попробует снова. Без этого индекса «неизменяемый поток» — просто договорённость, которую первая же гонка нарушит.

Как это выглядит в базе

Хранилище событий звучит внушительно, а в простейшем виде это одна таблица в PostgreSQL:

CREATE TABLE event_store (
    id          bigserial   PRIMARY KEY,
    stream_id   uuid        NOT NULL,          -- к какой сущности: заказ 42
    version     int         NOT NULL,          -- номер события в этом потоке: 1, 2, 3…
    type        text        NOT NULL,          -- 'OrderPaid'
    payload     jsonb       NOT NULL,          -- данные события
    metadata    jsonb       NOT NULL,          -- кто, откуда, идентификатор запроса
    occurred_at timestamptz NOT NULL DEFAULT now(),
    CONSTRAINT event_stream_version UNIQUE (stream_id, version)   -- вся защита от гонок здесь
);

CREATE INDEX idx_event_store_stream ON event_store (stream_id, version);

Что важно по полям:

  • stream_id и version вместе с уникальным индексом — это и есть весь механизм одновременной записи, о котором сказано выше: второй писатель с тем же номером получит отказ.
  • payload как документ, а не колонки: у разных типов событий разные поля, и схема таблицы не должна знать про каждый из них. Отсюда же и слабая схема, которая делает версионирование терпимым.
  • metadata отдельно от данных — кто совершил действие, идентификатор запроса и трассы, версия схемы события. Это то, что потом отвечает на вопросы поддержки, и его не смешивают с полезной нагрузкой.
  • id как сквозной номер нужен проекциям: потребитель помнит, до какого номера он дочитал, и продолжает с него.
  • Индекс по потоку — основной запрос («все события заказа 42 по порядку»).

И код применения — то, что превращает список событий в состояние:

public class Order {
    private OrderId id;
    private final Map<Integer, OrderLine> lines = new LinkedHashMap<>();
    private Money paid = Money.zero();
    private OrderStatus status = OrderStatus.NEW;
    private int version;

    // Восстановление: проиграть поток по порядку
    public static Order replay(List<StoredEvent> stream) {
        Order order = new Order();
        stream.forEach(order::apply);
        return order;
    }

    private void apply(StoredEvent stored) {
        switch (stored.event()) {
            case OrderCreated e -> { this.id = e.orderId(); this.status = NEW; }
            case ItemAdded    e -> lines.put(e.lineId(), new OrderLine(e.item(), e.price()));
            case ItemRemoved  e -> lines.remove(e.lineId());
            case OrderPaid    e -> { this.paid = e.amount(); this.status = PAID; }
        }
        this.version = stored.version();       // помним, на какой версии остановились
    }

    // Команда: проверить правила и вернуть НОВОЕ событие, ничего не записывая
    public OrderPaid pay(Money amount) {
        if (status != NEW) throw new IllegalTransition(status, "pay");
        if (!amount.equals(total())) throw new PartialPaymentNotAllowed(id);
        return new OrderPaid(id, amount);
    }

    public Money total() {
        return lines.values().stream().map(OrderLine::price).reduce(Money.zero(), Money::add);
    }
}
type Order struct {
	id      OrderID
	lines   map[int]OrderLine // состояние хранит карта: порядок позиций важен только витрине
	paid    Money
	status  OrderStatus
	version int
}

// Восстановление: проиграть поток по порядку
func Replay(stream []StoredEvent) *Order {
	order := &Order{lines: map[int]OrderLine{}, status: StatusNew}
	for _, stored := range stream {
		order.apply(stored)
	}
	return order
}

func (o *Order) apply(stored StoredEvent) {
	switch e := stored.Event.(type) {
	case OrderCreated:
		o.id, o.status = e.OrderID, StatusNew
	case ItemAdded:
		o.lines[e.LineID] = OrderLine{Item: e.Item, Price: e.Price}
	case ItemRemoved:
		delete(o.lines, e.LineID)
	case OrderPaid:
		o.paid, o.status = e.Amount, StatusPaid
	}
	o.version = stored.Version // помним, на какой версии остановились
}

// Команда: проверить правила и вернуть НОВОЕ событие, ничего не записывая
func (o *Order) Pay(amount Money) (OrderPaid, error) {
	if o.status != StatusNew {
		return OrderPaid{}, &IllegalTransition{From: o.status, Action: "pay"}
	}
	if amount != o.Total() {
		return OrderPaid{}, &PartialPaymentNotAllowed{ID: o.id}
	}
	return OrderPaid{OrderID: o.id, Amount: amount}, nil
}

func (o *Order) Total() Money {
	total := MoneyZero
	for _, l := range o.lines {
		total = total.Add(l.Price)
	}
	return total
}
export class Order {
  #id;
  #lines = new Map();
  #paid = Money.zero();
  #status = OrderStatus.NEW;
  #version = 0;

  // Восстановление: проиграть поток по порядку
  static replay(stream) {
    const order = new Order();
    for (const stored of stream) order.#apply(stored);
    return order;
  }

  #apply(stored) {
    const e = stored.event;
    switch (e.type) {
      case 'OrderCreated': this.#id = e.orderId; this.#status = OrderStatus.NEW; break;
      case 'ItemAdded':    this.#lines.set(e.lineId, new OrderLine(e.item, e.price)); break;
      case 'ItemRemoved':  this.#lines.delete(e.lineId); break;
      case 'OrderPaid':    this.#paid = e.amount; this.#status = OrderStatus.PAID; break;
    }
    this.#version = stored.version;        // помним, на какой версии остановились
  }

  // Команда: проверить правила и вернуть НОВОЕ событие, ничего не записывая
  pay(amount) {
    if (this.#status !== OrderStatus.NEW) throw new IllegalTransition(this.#status, 'pay');
    if (!amount.equals(this.total())) throw new PartialPaymentNotAllowed(this.#id);
    return { type: 'OrderPaid', orderId: this.#id, amount };
  }

  total() {
    return [...this.#lines.values()].reduce((sum, l) => sum.add(l.price), Money.zero());
  }
}
class Order:
    def __init__(self) -> None:
        self._id: OrderId | None = None
        self._lines: dict[int, OrderLine] = {}
        self._paid = Money.zero()
        self._status = OrderStatus.NEW
        self._version = 0

    # Восстановление: проиграть поток по порядку
    @classmethod
    def replay(cls, stream: list[StoredEvent]) -> "Order":
        order = cls()
        for stored in stream:
            order._apply(stored)
        return order

    def _apply(self, stored: StoredEvent) -> None:
        match stored.event:
            case OrderCreated(order_id=order_id):
                self._id, self._status = order_id, OrderStatus.NEW
            case ItemAdded(line_id=line_id, item=item, price=price):
                self._lines[line_id] = OrderLine(item, price)
            case ItemRemoved(line_id=line_id):
                del self._lines[line_id]
            case OrderPaid(amount=amount):
                self._paid, self._status = amount, OrderStatus.PAID
        self._version = stored.version        # помним, на какой версии остановились

    # Команда: проверить правила и вернуть НОВОЕ событие, ничего не записывая
    def pay(self, amount: Money) -> OrderPaid:
        if self._status is not OrderStatus.NEW:
            raise IllegalTransition(self._status, "pay")
        if amount != self.total():
            raise PartialPaymentNotAllowed(self._id)
        return OrderPaid(self._id, amount)

    def total(self) -> Money:
        return sum((line.price for line in self._lines.values()), Money.zero())

Два свойства этого кода — самое важное в реализации. Первое: apply ничего не проверяет. Он только меняет состояние, потому что проигрывает уже случившиеся факты: событие в журнале нельзя отвергнуть, оно уже произошло. Проверки живут в командах (pay), которые решают, можно ли добавить новое событие. Смешать это — частая и дорогая ошибка: проверка внутри apply однажды сделает старый поток непроигрываемым, потому что правила изменились.

Второе: команда возвращает событие, а не записывает его. Запись — работа обёртки:

@Transactional
public void handle(PayOrder command) {
    List<StoredEvent> stream = store.read(command.orderId());
    Order order = Order.replay(stream);

    OrderPaid event = order.pay(command.amount());        // правила

    store.append(command.orderId(), order.version() + 1, event);   // отказ при гонке
    outbox.save(event);                                    // наружу — через таблицу исходящих
}
func (h *PayOrderHandler) Handle(ctx context.Context, cmd PayOrder) error {
	return h.tx.RunInTx(ctx, func(ctx context.Context) error {
		stream, err := h.store.Read(ctx, cmd.OrderID)
		if err != nil {
			return err
		}
		order := Replay(stream)

		event, err := order.Pay(cmd.Amount) // правила
		if err != nil {
			return err
		}
		if err := h.store.Append(ctx, cmd.OrderID, order.Version()+1, event); err != nil { // отказ при гонке
			return err
		}
		return h.outbox.Save(ctx, event) // наружу — через таблицу исходящих
	})
}
handle(command) {
  return this.tx.runInTx(async () => {
    const stream = await this.store.read(command.orderId);
    const order = Order.replay(stream);

    const event = order.pay(command.amount);                            // правила

    await this.store.append(command.orderId, order.version + 1, event); // отказ при гонке
    await this.outbox.save(event);                                      // наружу — через таблицу исходящих
  });
}
def handle(self, command: PayOrder) -> None:
    with self._uow:
        stream = self._uow.store.read(command.order_id)
        order = Order.replay(stream)

        event = order.pay(command.amount)                                   # правила

        self._uow.store.append(command.order_id, order.version + 1, event)  # отказ при гонке
        self._uow.outbox.save(event)                                        # наружу — через таблицу исходящих
        self._uow.commit()

Строить самому или взять готовое

Таблица выше — полноценное хранилище событий, и для одной-двух сущностей этого достаточно. Что есть кроме неё:

Таблица в PostgreSQL. Плюсы: ничего нового в эксплуатации, транзакция общая с бизнес-данными и с таблицей исходящих сообщений, запросы обычным SQL, резервные копии как у всего остального. Минусы: проекции, снапшоты, подписку на поток и повторное проигрывание надо написать самому (это несколько сотен строк, не тысячи). Разумный выбор по умолчанию, особенно когда событийная модель применяется точечно.

Библиотека поверх своей базы (в мире Java это Axon Framework, в .NET и Node есть аналоги). Даёт готовыми: агрегаты с проигрыванием, снапшоты, проекции, обработку команд, иногда и распределённую доставку. Цена: фреймворк начинает определять устройство кода, и уйти от него потом дорого. Оправдан, когда событийная модель — основа всей системы, а не одна сущность.

Специализированное хранилище (EventStoreDB и подобные). Даёт подписки на потоки, проекции на стороне сервера, оптимизированное чтение, категории потоков. Цена: ещё одна система в эксплуатации, отдельные резервные копии, отсутствие общей транзакции с вашей базой — а значит, таблица исходящих и идемпотентность становятся обязательными. Берут при больших объёмах событий или когда нужны серверные подписки.

Чего делать не стоит: использовать Kafka как хранилище событий «потому что там тоже журнал». Kafka — транспорт с ограниченным сроком хранения и без чтения по ключу: получить «все события заказа 42» из темы нельзя, а восстановление состояния одной сущности превращается в перебор. Kafka рядом с хранилищем событий — нормально и часто (доставка проекциям и другим сервисам); Kafka вместо него — нет.

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

Как из событий собирается состояние

Чтобы узнать текущее состояние заказа, мы берём все его события по порядку и применяем одно за другим к пустой заготовке. Этот процесс называется replay (проигрывание):

старт заказ пустой OrderCreated заказ 42, покупатель Аня, сумма 0 ItemAdded книга 500 → сумма 500 ItemAdded ручка 100 → сумма 600 ItemRemoved ручка убрана → сумма 500 OrderPaid оплачен, сумма 500 итог заказ 42, оплачен, книга, сумма 500

Состояние нигде не хранится — оно получается прокруткой событий с начала. Отсюда и главное свойство: видно не только «сколько сейчас», но и как к этому пришли.

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

Проекции и связь с CQRS

Проигрывать весь поток каждый раз, когда кто-то открывает список заказов, — дорого и медленно. Поэтому Event Sourcing почти всегда идёт в паре с CQRS — разделением на запись и чтение.

Вопросов два: кто дописывает факты и кто собирает из них удобную для чтения картину. Работает это так:

  • Командная сторона (запись) принимает команду, проверяет правила и дописывает в event store новое событие. Она работает только с потоком фактов.
  • Запросная сторона (чтение) заранее строит удобные для чтения таблицы — проекции (их ещё называют read-model). Проекция подписана на поток событий: пришло OrderPaid — она обновила строку в таблице orders_view, где лежит уже готовое состояние.

Отдельный вопрос — как событие доезжает до проекции. Обычно его вычитывают из event store отдельным процессом или получают через очередь, и в обоих случаях доставка работает по правилу «хотя бы один раз»: одно и то же событие может прийти дважды. Для проекции это опасно — повторно применённое ItemAdded добавит позицию во второй раз. Поэтому проекция запоминает, до какого номера версии она обработала поток, и всё, что не больше этого номера, молча пропускает. Разбор с примерами — в статье «Синхронизация через события» (Java, Go, Node, Python).

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

команда журнал событий дописывает факт список заказов экран клиента отчёт по выручке финансы аналитика собрана позже

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

Снапшоты: чтобы не проигрывать всё с нуля

У долгоживущей сущности событий могут быть десятки тысяч. Проигрывать их все, чтобы восстановить состояние, — медленно. Решение — снапшот (snapshot): периодически сохранять «слепок» состояния на определённой позиции в потоке.

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

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

событие 1 заказ создан ... ещё 998 событий событие 1000 снимок состояния событие 1001 позиция добавлена событие 1002 заказ оплачен восстановление снимок и два события

Снимок на позиции 1000 отрезает предысторию: при восстановлении читают его и только события после него, а сам снимок всегда можно выбросить и собрать заново.

С одной оговоркой, и она важнее, чем кажется: пересобрать получится, только если логика применения событий не поменялась. Если год назад OrderPaid считался иначе, чем считается сегодня, то старый снапшот и заново собранное состояние разойдутся. Поэтому снапшоты и версионирование событий — про одно и то же: как только меняется правило применения, снапшоты, снятые по старому правилу, надо выбрасывать и пересчитывать.

Запросы, которые не сводятся к одному потоку

Главное ограничение подхода, и его стоит произнести прямо: из журнала событий нельзя спросить «все заказы за март дороже 1000». Журнал устроен для чтения одного потока по идентификатору; любой вопрос про множество сущностей, фильтр или сумму в нём не выражается.

Почему так: состояние заказа — это результат проигрывания его событий, и оно не лежит в базе в готовом виде. Чтобы отфильтровать заказы по сумме, пришлось бы проиграть все потоки и посчитать сумму у каждого. На тысяче заказов это работает; на миллионе — нет.

Ответ один, и он не обходится: такие запросы обслуживает проекция — обычная таблица с текущим состоянием, которую наполняют обработчики событий. Заказы дороже 1000 за март читаются из order_summary простым запросом с индексом, как в любом сервисе.

Практические следствия, которые стоит принять сразу:

  • Проекция — не оптимизация, а обязательная часть архитектуры. Без неё у системы нет способа отвечать на большинство вопросов интерфейса. Тот, кто внедряет событийную модель «сначала журнал, проекции потом», обнаруживает, что «потом» — это первый же экран списка.
  • Проекций обычно несколько. Список заказов, витрина для поддержки, отчётные агрегаты, поиск — у каждого свои поля и свои индексы. Это нормально и дёшево: все они выводятся из одного журнала.
  • Данные в проекции отстают. На доли секунды или секунды — и это то самое свойство, которое надо закладывать в интерфейс (см. минусы ниже).
  • Проекцию можно пересобрать — и это её лучшее свойство: сломали, изменили формат, добавили поле — очистили и проиграли журнал заново. Значит, у системы обязана быть процедура перестроения проекции, и её надо проверить до того, как понадобится: на живом объёме это часы, и об этом лучше знать заранее.
  • Уникальность и ограничения не работают через журнал. «Не может быть двух активных подписок у одного клиента» — правило, которое нельзя проверить в журнале одним запросом; его проверяют либо внутри одного потока (если сущность одна), либо через проекцию с уникальным индексом (и тогда с отложенностью приходится считаться), либо ограничением на бизнес-процесс. Это одна из причин, по которой границы потоков проектируют так, чтобы правила укладывались внутрь одного потока.

И честная формулировка ограничения целиком: журнал отвечает на «что случилось с этой сущностью», проекция — на все остальные вопросы. Система без проекций — не событийная, а нерабочая.

Плюсы: полный аудит и «машина времени»

Главная ценность в том, что вы никогда ничего не теряете.

  • Полный аудит из коробки. Вопрос «кто, что и когда менял» отвечается сам собой — история изменений и есть ваша модель данных, а не отдельный лог, который забыли вести.
  • «Машина времени». Можно восстановить состояние на любой момент прошлого: просто проиграть события до нужной точки. Удобно для разбора инцидентов — «а как выглядел заказ вчера в 15:00?».
  • Отладка и воспроизведение. Баг проявился в проде? Можно проиграть те же события у себя и увидеть ровно то же состояние.
  • Новые отчёты задним числом. Понадобилась метрика, о которой год назад не думали, — строите новую проекцию и проигрываете всю историю.

Минусы: почему это дорого и нужно редко

Всё это не бесплатно, и честно сказать: Event Sourcing нужен куда реже, чем кажется на волне интереса к нему.

  • Сложность. Разработчику приходится думать не «обнови поле», а «какое событие произошло и как оно применяется». Порог входа выше, отладка непривычнее.
  • Эволюция схемы событий. События хранятся вечно, а формат со временем меняется. Событие OrderPaid пятилетней давности должно проигрываться и сегодня. Приходится версионировать события и уметь читать старые форматы — это постоянная работа.
  • Отложенная согласованность проекций. Между записью события и обновлением read-model проходит время. Пользователь оплатил заказ, а в списке пару секунд ещё «не оплачен». Это eventual consistency, и её надо закладывать в интерфейс, а не бороться с ней.
  • Инфраструктура. Event store, механизм проекций, снапшоты, доставка событий — всё это надо построить и поддерживать. Часто рядом появляется Kafka как транспорт для событий, и вместе с ней — свои заботы о надёжности доставки.
  • Ошибку нельзя исправить правкой строки. Самый болезненный пункт в эксплуатации, и он меняет практику сильнее порога входа: обычно ошибку в данных чинят UPDATE, а здесь неверное событие уже записано и переписать его нельзя (уникальный номер в потоке и append-only — это не соглашение, а устройство). Исправление возможно только новым, компенсирующим событием.

Что делать, когда в журнал попал неверный факт

Сценарий обычный: ошибка в обработчике команды записала OrderPaid с суммой в копейках вместо рублей, и это уехало в проекции. Порядок действий.

1. Остановить источник. Сначала починить код, потом разбираться с данными: пока ошибка жива, поток продолжает наполняться неверными событиями.

2. Записать компенсирующее событие. Не «исправить сумму», а зафиксировать новый факт:

OrderPaid          { orderId: 42, amount: 50000 }      // ошибочный факт, остаётся навсегда
PaymentCorrected   { orderId: 42, wrongAmount: 50000,  // новый факт
                     correctAmount: 500, reason: "bug-2841", correctedBy: "support:ivanov" }

Обработчик состояния учитывает оба: сумма берётся с учётом коррекции. Именно так работает бухгалтерия — исправляющей проводкой, а не подчисткой, — и событийная модель ведёт себя так же по построению.

3. Научить проигрывание понимать коррекцию. Событие коррекции — часть модели, а не служебная запись: у него есть тип, правила применения и тесты. Это работа, и её нельзя избежать: проекции и снапшоты должны считать так же, как состояние.

4. Пересобрать проекции. После изменения логики применения — очистить и проиграть журнал заново, иначе витрина останется с неверными числами. Снапшоты, снятые по старому правилу, тоже выбрасывают.

Почему нельзя «просто подправить журнал». Формально доступ к базе есть, и запрос UPDATE event_store SET payload = ... выполнится. Что будет дальше: разойдутся снапшоты (они считались по старому значению), разойдутся проекции (они уже применили старое событие), исчезнет соответствие между журналом и тем, что видели пользователи и внешние системы, — и, главное, пропадёт свойство, ради которого всё это строилось: журнал больше не история, а редактируемая таблица. Если правка журнала допустима, событийная модель не нужна: обычный UPDATE дешевле.

Единственное исключение — требования закона об удалении персональных данных, и там ответ другой: не правка, а невозможность прочитать (шифрование персональных полей ключом на субъекта и уничтожение ключа). Разбор — в разделе «Глубже» ниже.

Что из этого следует для эксплуатации. Три вещи, которые стоит завести до внедрения: набор компенсирующих событий на предсказуемые случаи (коррекция суммы, отмена ошибочной операции, перенос на другой объект); процедура перестроения проекций с известным временем; и строгая проверка на входе — раз ошибку нельзя убрать, проверять надо до записи, а не после, и это меняет распределение усилий по сравнению с обычной моделью.

Поэтому Event Sourcing оправдан там, где история — часть требований, а не приятный бонус: финансы и платежи, где аудит обязателен по закону; сложный домен, где важна причинно-следственная цепочка решений; системы, где регулярно нужны разборы «как мы сюда пришли». Для обычного CRUD, где достаточно текущего состояния, привычные UPDATE проще, дешевле и надёжнее — и это нормальный выбор по умолчанию.

Где это применяется

Event Sourcing — не «продвинутый способ хранить данные», а инструмент под конкретную задачу: когда сама история изменений имеет ценность.

  • Платежи и бухгалтерия. Поток операций — естественная модель, а аудит требуется по регламенту. Здесь событийная модель ложится идеально.
  • Сложные бизнес-процессы. Заказы, страховые дела, заявки с длинным жизненным циклом и множеством переходов — когда важно видеть всю цепочку, а не только финал.
  • Системы с требованием «машины времени». Там, где регулярно нужно восстанавливать состояние на прошлый момент или разбирать инциденты по шагам.

Где спотыкаются начинающие:

  • Берут Event Sourcing на весь проект «чтобы по-взрослому», хотя 90% сущностей — обычный CRUD. Событийную модель применяют точечно, к тем частям домена, где история реально нужна, а не ко всему подряд.
  • Путают Event Sourcing и просто отправку событий в очередь. Публиковать OrderPaid в Kafka — это интеграция между сервисами. Event Sourcing — это когда поток событий и есть ваш источник правды, из которого вычисляется состояние. Это разные вещи, хотя события звучат похоже.
  • Забывают про версионирование событий и через год не могут проиграть старый поток, потому что формат события изменился. Совместимость со старыми событиями надо планировать с самого начала.
  • Ждут от read-model мгновенной согласованности и удивляются, почему только что оплаченный заказ секунду висит «неоплаченным». Отложенность проекций — это свойство подхода, а не баг.
Дополнительно: при первом чтении можно пропустить

Глубже: версии событий в журнале: upcasting и правило «только добавлять»расширенное

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

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

Номер версии в конверте. У каждого сохранённого события есть тип и версия схемы: OrderPaid версии 2. Семантическое изменение (сумма в копейках вместо рублей) это новая версия, а не правка старой; старые события в журнале не трогают никогда, они факт.

Подъём при чтении. Читатель журнала знает текущую версию каждого типа и держит цепочку преобразователей: из версии 1 в 2, из 2 в 3. Событие поднимается до текущей до того, как попадёт в код восстановления состояния, и сам агрегат знает только последнюю версию. Преобразователи это код, который живёт, пока в журнале есть старые события, то есть вечно, поэтому число версий держат малым, а преобразователи покрывают тестами на примерах реальных старых событий из журнала.

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

Снапшоты при смене версии сбрасывают и строят заново, потому что они производное, а проекции перестраивают из журнала, как описано выше; это и есть проверка, что подъём версий работает.

Глубже: персональные данные в неизменяемом журналерасширенное

Журнал, который никогда не редактируют, и закон, который требует удалить всё о человеке по запросу, противоречат друг другу напрямую, а статья рекомендует подход для финансов, где персональных данных много. Стандартный ответ называют криптографическим уничтожением.

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

Цена: хранилище ключей становится критичным компонентом, ключ на субъекта это миллионы ключей, каждое чтение события с персональными данными это расшифровка, и перестроение проекции идёт медленнее. Поэтому персональных полей в событиях держат мало: событие ссылается на клиента по идентификатору, а имя и адрес живут в обычной таблице, которую можно обновить и удалить; в журнал попадает только то, что нужно для восстановления состояния, а не для отображения.

Второе, о чём спрашивают при проверке: срок хранения. Журнал «вечный» по замыслу, а закон требует хранить не дольше цели; для финансовых событий цель и срок задаёт отраслевое регулирование (годы), и это записывают в документ обработки данных до того, как журнал заработал, о чём статья про 152-ФЗ.

Коротко

  • Event Sourcing хранит не состояние, а поток событий, из которого состояние восстанавливают; снапшоты ускоряют восстановление, проекции дают чтение.
  • Плюсы: полный аудит, машина времени, новые проекции из старых событий; минусы: сложность, эволюция схемы, отложенные проекции, инфраструктура.
  • Берут точечно там, где история это требование (платежи, сложные процессы), а не на весь проект.
  • Версии событий: только добавлять поля, номер версии в конверте, подъём при чтении цепочкой преобразователей, новый тип для нового смысла, журнал не редактируют.
  • Персональные данные шифруют ключом на субъекта и удаляют ключом; в журнал кладут минимум личных полей, остальное по ссылке в обычную таблицу.
  • Хранилище событий в простейшем виде — таблица с потоком, номером версии, типом, документом и метаданными: уникальный индекс по потоку и номеру и есть вся защита от одновременной записи.
  • Применение события ничего не проверяет (факт уже случился), проверки живут в командах, а команда возвращает событие, не записывая его; запись и публикация — работа обёртки в одной транзакции.
  • Выбор хранилища: таблица в своей базе по умолчанию, библиотека когда событийная модель — основа системы, специализированное хранилище при объёмах и серверных подписках; Kafka вместо хранилища не годится.
  • Из журнала нельзя спросить «все заказы за март дороже 1000»: такие вопросы обслуживает проекция, она обязательная часть архитектуры, её несколько, она отстаёт и её надо уметь пересобрать.
  • Неверное событие нельзя переписать: только компенсирующее событие, изменение логики применения и пересборка проекций; поэтому заводят набор компенсаций, процедуру перестроения и строгую проверку на входе.

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