В CQRS приложение делится на две части: запись (command side) и чтение (query side). Эта статья про первую — как команды меняют состояние, кто это делает и по каким правилам.
Что такое команда
Представьте кассира, который нажимает «Подтвердить заказ». Это намерение изменить состояние системы. В коде оно оформляется как команда — объект с данными, который несёт это намерение.
public record ConfirmOrderCommand(
Long orderId,
String idempotencyKey
) implements UseCaseCommand<OrderId> {}
Несколько важных деталей:
- Команда — это record (или final-класс): иммутабельный, с equals/hashCode по полям. Удобно логировать, передавать, тестировать.
- В конструкторе нет логики. Команда — это данные, не вычисления. Контроллер заполняет её из запроса и передаёт дальше.
- Параметр
<OrderId>— что вернёт обработчик. Об этом ниже. idempotencyKey— ключ для денежных операций, берётся из заголовкаIdempotency-Key.
Одна команда — один агрегат
Это ключевое правило command side. Одна транзакция меняет один агрегат.
Почему это важно? Посмотрим на проблемный пример:
// Частая ошибка — два агрегата в одной транзакции
@Transactional
public OrderId handle(CreateOrderCommand cmd) {
Customer customer = customerRepository.findById(cmd.customerId(), FOR_UPDATE)
.orElseThrow(() -> new CustomerNotFoundException(cmd.customerId()));
customer.incrementOrderCount(); // меняем Customer
customerRepository.save(customer);
Order order = orderFactory.createFor(customer, cmd.items());
orderRepository.save(order); // меняем Order
return order.id();
}
Здесь транзакция держит блокировки сразу на двух строках БД. При конкурентной нагрузке это означает:
- больше ожидания (другие операции с Customer или Order стоят в очереди);
- риск взаимной блокировки (deadlock), если другая транзакция возьмёт блокировки в другом порядке;
- при откате — оба изменения теряются вместе.
Правильно: один агрегат в транзакции, остальные обновляются асинхронно через события:
// Правильно — меняем только Order
@Transactional
public OrderId handle(CreateOrderCommand cmd) {
Order order = orderFactory.createFor(cmd.customerId(), cmd.items());
orderRepository.save(order);
// OrderCreated-событие зарегистрировано внутри агрегата;
// relay опубликует его, Customer-сервис обновит счётчик сам
return order.id();
}
Если бизнес-логика требует менять два агрегата — это либо saga (цепочка команд с компенсациями), либо сигнал, что границы агрегатов нарезаны неверно.
Что такое сага, в двух словах и одном примере
Слово встречается в статье трижды, и без содержания оно бесполезно. Сага — это последовательность локальных транзакций в разных агрегатах или сервисах, у каждой из которых есть компенсирующее действие. Общей транзакции нет: вместо откатов — отмены уже сделанного.
Оформление заказа через три агрегата:
| Шаг | Прямое действие | Компенсация |
|---|---|---|
| 1 | Заказ переведён в «ожидает оплаты» | Заказ отменён |
| 2 | Товар зарезервирован на складе | Резерв снят |
| 3 | Деньги списаны | Деньги возвращены |
Отказ на третьем шаге запускает компенсации в обратном порядке: снять резерв, отменить заказ. Именно «в обратном» — потому что поздние шаги могли опираться на результат ранних.
Кто её запускает — главный вопрос, и ответов два. Координатор (оркестрация): отдельный компонент знает всю последовательность, вызывает шаги по очереди, ловит отказы и запускает компенсации. Видно весь процесс в одном месте, легко отлаживать, легко добавить шаг; цена — компонент, который знает про всех. Без координатора (координация участниками): каждый участник слушает событие предыдущего и реагирует; при отказе публикует событие отказа, а остальные сами снимают своё. Участники независимы; цена — процесс не виден целиком, а зависшая сага никем не замечена. Разбор обоих способов и границы выбора — в статье про движки процессов.
Три свойства, без которых сага не работает. Первое: каждый шаг идемпотентен — повтор не создаёт второй резерв и не списывает дважды. Второе: у саги есть состояние — где она сейчас, какие шаги прошли, что придётся компенсировать; это таблица, а не память процесса. Третье: есть сторож — задача, которая находит саги, застрявшие дольше срока, и доводит их (компенсацией или повтором); без него подвисшая сага висит навсегда.
Что делать, если компенсация сама упала. Самый неприятный вопрос темы, и ответ честный: компенсацию повторяют, потому что она обязана выполниться, а после нескольких неудач — ручной разбор. Значит, у компенсаций есть свой счётчик попыток и своя тревога: упавшая компенсация оставляет систему в противоречивом состоянии (деньги списаны, резерв снят, заказ отменён), и это тот случай, когда человек нужен. Отсюда практическое следствие: компенсация должна быть проще прямого действия — «вернуть деньги» надёжнее, чем «откатить всю сделку», и проектируют её так же тщательно, как основной путь.
И когда сагу не заводят. Если два агрегата меняются вместе всегда и расхождение между ними недопустимо ни на секунду — это признак того, что граница агрегатов проведена неверно, и правильный ответ не сага, а слияние в один агрегат. Сага — для процессов, где между шагами допустимо промежуточное состояние.
Путь команды от HTTP до таблицы исходящих: транзакция открывается в обработчике и закрывается после сохранения, внешние вызовы остаются за её границами.
Структура command-handler-а
Обработчик команды (command-handler) — это класс с одним методом handle. Его работа всегда состоит из одних и тех же шагов:
@Component
@RequiredArgsConstructor
class ConfirmOrderHandler implements UseCaseHandler<ConfirmOrderCommand, OrderId> {
private final OrderRepository orderRepository;
private final Clock clock;
@Override
@Transactional
public OrderId handle(ConfirmOrderCommand cmd) {
// 1. Загрузить агрегат с блокировкой
Order order = orderRepository.findById(
new OrderId(cmd.orderId()),
SelectMode.FOR_UPDATE)
.orElseThrow(() -> new OrderNotFoundException(cmd.orderId()));
// 2. Вызвать доменный метод — он проверяет инварианты
order.confirm(clock.instant());
// 3. Сохранить (событие OrderConfirmed уже внутри агрегата)
orderRepository.save(order);
// 4. Вернуть минимум
return order.id();
}
}
Разберём каждый шаг.
Блокировка при загрузке (FOR UPDATE)
SelectMode.FOR_UPDATE — это пессимистическая блокировка: пока текущая транзакция не завершится, другая транзакция не сможет взять ту же строку под FOR UPDATE и не сможет её изменить — она встанет и будет ждать. Обычному чтению это не мешает: SELECT без блокировки прочитает прежнюю версию строки и не остановится.
Зачем это нужно? Без блокировки возможна классическая гонка: две транзакции одновременно читают Order.status = NEW, обе решают подтвердить заказ, обе записывают — одно из изменений теряется.
Пессимистическая блокировка — не единственный способ от этого защититься. Второй — оптимистическая блокировка: у агрегата есть колонка version, при сохранении она проверяется и увеличивается, и если её успели изменить — сохранение падает, а команда повторяется с начала. Разница по смыслу: FOR UPDATE не пускает конкурента сразу, оптимистическая пускает и ловит на выходе. Когда за одну и ту же строку дерутся часто, дешевле первый вариант; когда конфликты редки, а строка горячая (все команды идут в один заказ), очередь на FOR UPDATE сама становится узким местом — и лучше второй.
Инварианты — внутри агрегата
Обработчик не проверяет состояние сам. Он вызывает доменный метод, который сам знает, что допустимо:
public final class Order extends AggregateRoot<OrderId> {
public void confirm(Instant now) {
if (this.status != OrderStatus.NEW) {
throw new OrderAlreadyConfirmedException(this.id, this.status);
}
if (this.items.isEmpty()) {
throw new EmptyOrderException(this.id);
}
this.status = OrderStatus.CONFIRMED;
registerEvent(new OrderConfirmed(this.id, now));
}
}
Два места здесь стоит пояснить отдельно.
AggregateRoot<OrderId> — базовый класс из библиотеки ddd-building-blocks. Он даёт агрегату идентификатор и складывает зарегистрированные события в список: registerEvent не отправляет ничего никуда, он просто кладёт событие «в карман» агрегата, откуда его заберёт репозиторий при сохранении.
Время приходит параметром, а не берётся через Instant.now() внутри агрегата. Системные часы — такой же внешний мир, как база: если домен читает их сам, проверить поведение «через тридцать дней» можно только отладочным методом. Значение now берёт обработчик у бина Clock, а в тесте вместо него стоит Clock.fixed(...).
Правило простое: логика «можно ли это сделать» живёт в агрегате, обработчик только оркестрирует поток (загрузить → вызвать → сохранить).
Событие регистрируется внутри агрегата
order.confirm(now) вызывает registerEvent(new OrderConfirmed(...)). При save репозиторий публикует это событие в outbox. Relay доставит его подписчикам (например, сервису синхронизации read-модели). Обработчик всё это не знает и не координирует вручную.
Что возвращает команда
Команда возвращает минимум: идентификатор, статус или пустой результат. Не полный read-DTO.
// Возвращает только id
public record CreateOrderCommand(...) implements UseCaseCommand<OrderId> {}
// Не возвращает ничего значимого
public record CancelOrderCommand(...) implements UseCaseCommand<UseCaseEmptyResult> {}
// Частая ошибка — возвращает полный read-DTO
public record CreateOrderCommand(...) implements UseCaseCommand<OrderSummaryJson> {}
// Частая ошибка — возвращает сам агрегат
public record CreateOrderCommand(...) implements UseCaseCommand<Order> {}
Отдавать наружу агрегат нельзя по той же причине, по которой его не отдаёт и query-обработчик: получив Order, контроллер сможет вызвать order.confirm() напрямую, в обход обработчика и всех проверок. Агрегат остаётся внутри, наружу уходит идентификатор.
Почему нельзя возвращать полный read-DTO? Потому что это смешивает две обязанности: запись и чтение. Command-handler занимается записью. Если клиенту нужна полная проекция — он делает отдельный GET-запрос после записи. Контракт: POST /orders возвращает 201 с Location: /orders/{id}, а клиент при необходимости читает через GET /orders/{id}.
Четыре возможных результата команды: сверху то, что вернуть можно, снизу то, что ломает разделение чтения и записи.
Решение зависит от чужого агрегата
Правило «читать состояние для решения — только из своего агрегата» оставляет открытым главный практический случай: нужное поле физически принадлежит другому агрегату, и втащить его в свой нельзя. Например, «нельзя подтвердить заказ, если у покупателя просрочена оплата предыдущего» — просрочка живёт в платежах, а не в заказе. Что делать, по возрастанию цены и надёжности.
1. Передать решение в команду, а не данные. Самый дешёвый и самый недооценённый приём: проверку делает тот, кто её умеет, до команды, и в команду приходит уже готовый ответ.
// Обработчик HTTP-запроса или сервис-фасад — до транзакции команды
PaymentEligibility eligibility = paymentsClient.eligibilityFor(customerId); // чужой контекст
confirmOrder.handle(new ConfirmOrderCommand(orderId, eligibility.allowed(), key));
Агрегат при этом принимает решение сам, но по переданному факту, а не по данным чужого агрегата: order.confirm(eligibilityAllowed). Правило остаётся соблюдённым, транзакция короткая, и внешний вызов — вне неё. Ограничение: факт может устареть между проверкой и выполнением — и вот здесь нужен следующий пункт.
2. Скопировать нужный факт в свой агрегат через событие. Если решение требуется часто, чужой факт переносят к себе: платежи публикуют «у покупателя просрочка», заказы держат у себя поле customerHasOverdue и обновляют его по событию. Дальше решение принимается локально и мгновенно. Это не дублирование данных, а проекция чужого факта в свою модель, и она стоит ровно того, чего стоит любая проекция: отставание на доли секунды и необходимость уметь пересобрать. Признак, что приём подходит: факт меняется редко, а читается часто.
3. Принять решение оптимистично и компенсировать. Когда проверить заранее нельзя (резерв на складе, лимит, баланс), команда выполняется, а чужая сторона подтверждает или отклоняет — то есть получается сага из предыдущего раздела. Пользователь видит промежуточное состояние («ожидает подтверждения»), и это честнее, чем блокировать два агрегата в одной транзакции.
4. Признать, что границы неверные. Если решение всегда требует чужого состояния и промежуточное состояние недопустимо, значит, два агрегата — на самом деле один. Это редкий, но реальный случай; признак — каждое изменение одного требует одновременного изменения другого.
Чего делать не стоит. Читать чужую таблицу напрямую из своего обработчика (связывает схемы и обходит владельца данных), держать чужой агрегат в транзакции «только для чтения» (это всё равно блокировка и всё равно связанность), и делать HTTP-вызов в чужой сервис внутри транзакции — про последнее подробно в разделе «Глубже» ниже.
Как выбрать: редко меняется и часто нужно — копировать событием; нужно однократно и точно — передавать решением до команды; нельзя знать заранее — сага. Четвёртый вариант оставляют на случай, когда первые три не помогли.
Валидация: где что проверяется
Валидация происходит в двух разных местах, и их нельзя путать.
Контракт запроса — проверяется на контроллере через Jakarta Validation:
public record ConfirmOrderRequest(
@NotNull Long orderId,
@NotBlank @Size(max = 64) String idempotencyKey
) {}
Это технические ограничения: поле не пустое, размер в пределах, формат правильный.
Бизнес-инвариант — проверяется в методе агрегата и бросает доменное исключение:
public void confirm(Instant now) {
if (this.status != OrderStatus.NEW) {
throw new OrderAlreadyConfirmedException(this.id, this.status);
}
// ...
}
Это бизнес-правило: заказ можно подтвердить только если он в статусе NEW. Такая логика принадлежит домену, не HTTP-слою.
Частые ошибки
Отдельный SELECT «прочитать и решить»
// Так делать не нужно
@Transactional
public OrderId handle(ConfirmOrderCommand cmd) {
boolean hasPayment = paymentRepository.existsByOrderId(cmd.orderId()); // отдельный read
if (!hasPayment) throw new PaymentRequiredException(cmd.orderId());
Order order = orderRepository.findById(new OrderId(cmd.orderId()), FOR_UPDATE)
.orElseThrow(...);
order.confirm(clock.instant());
orderRepository.save(order);
return order.id();
}
Проблема двойная: между первым и вторым запросом данные могут измениться (гонка без блокировки), и read-логика просочилась в write-handler. Если paymentStatus важен для подтверждения заказа — он должен быть полем агрегата Order. Тогда order.confirm(now) сам проверит его.
Два агрегата в одной транзакции
Уже разобрали выше. Перевод денег между двумя счетами — классический пример, где нужна saga:
// Частая ошибка
@Transactional
public void handle(TransferMoneyCommand cmd) {
Account from = accountRepository.findById(cmd.fromId(), FOR_UPDATE).orElseThrow(...);
Account to = accountRepository.findById(cmd.toId(), FOR_UPDATE).orElseThrow(...);
from.debit(cmd.amount());
to.credit(cmd.amount());
accountRepository.save(from);
accountRepository.save(to);
}
// Правильно — saga из двух локальных команд:
// 1. DebitAccount → AccountDebited
// 2. CreditAccount ← orchestrator по событию
// При ошибке — компенсирующая команда CreditAccount (возврат)
Глубже: версия агрегата: откуда берётся и что видит проигравшийрасширенное
Идемпотентность проекций в статьях про синхронизацию держится на колонке version и поле aggregateVersion в событии, а command side учит только FOR UPDATE. Откуда версия берётся и кто её увеличивает, стоит показать здесь, потому что это часть обработчика.
Версия это счётчик на корне агрегата, который растёт на единицу при каждой успешной команде. Обработчик загружает агрегат вместе с версией, выполняет команду, и репозиторий пишет с условием: UPDATE orders SET ..., version = 8 WHERE id = ? AND version = 7. Обновилась одна строка, команда прошла и событие уходит с aggregateVersion = 8; обновилось ноль строк, кто-то успел раньше, и репозиторий бросает исключение оптимистичной блокировки. Это дешевле FOR UPDATE: никто никого не ждёт, а конфликты редки.
Проигравший видит исключение, и что с ним делать, решает обработчик, а не клиент. Команда идемпотентна по смыслу (подтвердить заказ): перечитать агрегат и повторить, до трёх раз, и повтор в обработчике прячет конфликт от клиента. Команда зависит от того, что клиент видел (изменить адрес после того, как другой сотрудник его тоже менял): вернуть клиенту 409 с текущей версией, чтобы он перечитал и решил. Команда с деньгами: повтор только с ключом идемпотентности из следующего раздела, иначе повторная попытка может списать дважды.
FOR UPDATE остаётся для случаев, когда конфликты постоянны (счётчик, который меняют сотни раз в секунду) или когда нужно посчитать что-то по нескольким строкам до записи; там оптимистичная блокировка даст поток повторов. Механику обеих и SelectMode в репозитории разбирает статья про реализацию автомата, а размер агрегата как причину постоянных конфликтов статья про тактические паттерны DDD.
Глубже: ключ идемпотентности команды: где хранится и что вернуть на повторрасширенное
Поле idempotencyKey появилось в команде выше и больше не упоминалось, а для денежных команд это центральная деталь: клиент отправил «оплатить», ответ потерялся, клиент повторил, и второй платёж это инцидент.
Ключ генерирует клиент (UUID на попытку операции, одинаковый при повторах) и присылает в заголовке Idempotency-Key; для внутренних команд ключ естественный, например «возврат по заказу 42». Обработчик первым действием в транзакции вставляет строку в таблицу command_idempotency(key, actor, request_hash, status, response, created_at) с уникальным ключом по паре «ключ, актор». Вставка прошла, значит команда новая, и по завершении в ту же строку пишут результат (статус и тело ответа) той же транзакцией. Вставка упала на уникальности, значит повтор, и дальше три исхода: строка в состоянии «выполнено» и хеш запроса совпал, вернуть сохранённый ответ с тем же кодом, как в первый раз; хеш не совпал, тот же ключ с другим телом, ответ 422; строка «в работе», первая попытка ещё идёт, ответ 409, чтобы клиент подождал и повторил.
Срок жизни строк не меньше суток и обычно неделя, потому что клиент может повторить через час после сбоя; чистка по возрасту. Ключ привязан к актору, иначе чужой ключ вернёт чужой ответ. И запись ключа не отделяют от команды: ключ в кэше Redis, а команда в базе, это окно, где повтор проходит; всё в одной транзакции той же базы, о чём статьи про идемпотентность в корректном завершении и про заголовки REST.
Глубже: границы транзакции в обработчикерасширенное
@Transactional стоит на обработчиках во всех примерах, и стоит сказать, что он делает в четырёх ситуациях, где интуиция обманывает.
Внешний вызов внутри транзакции. Обработчик в транзакции зовёт платёжный сервис по HTTP на три секунды: всё это время соединение с базой занято, строки заблокированы, а если вызов прошёл, а транзакция откатилась, деньги списаны, заказа нет. Правило: внешние вызовы до транзакции (проверить, можно ли) или после через outbox (сделать по событию); внутри транзакции только база.
Исключение внутри обработчика. Транзакция откатывается только на непроверяемых исключениях; проверяемое исключение по умолчанию фиксирует всё, что успело записаться, и это старая ловушка Spring, которую закрывают rollbackFor = Exception.class или отказом от проверяемых исключений в доменном коде. Исключение, пойманное внутри и не проброшенное, откат не вызывает, но если оно прошло через вложенный @Transactional, транзакция уже помечена на откат, и фиксация упадёт с «transaction silently marked rollback-only».
Событие при откате. Событие, опубликованное через ApplicationEventPublisher внутри транзакции, слушатель с @TransactionalEventListener получит только после фиксации, и при откате не получит вовсе, что правильно; обычный @EventListener получит его сразу, и при откате событие уже разлетелось. Событие в outbox откатывается вместе с транзакцией само, потому что это строка в той же базе.
Вызов из того же класса. @Transactional работает через прокси, и вызов метода изнутри того же класса идёт мимо прокси: транзакции нет, хотя аннотация стоит. Обработчик команды это один публичный метод handle, и вложенные транзакционные методы в нём не заводят.
Остальное из статьи про @Transactional в разделе Spring: readOnly для запросов, таймаут на транзакцию, короткие транзакции без ожидания пользователя.
Коротко
- Команда — record с данными и маркером
UseCaseCommand<R>. Без логики в конструкторе. Одна команда меняет один агрегат. Два агрегата — это saga. - От потери обновлений при конкурентных запросах защищаются двумя способами:
FOR UPDATEпри загрузке или колонкаversionс повтором команды при конфликте. Совсем без защиты нельзя. Инварианты проверяет агрегат, не обработчик. Handler оркеструет: загрузить → вызвать метод → сохранить. - Событие регистрируется внутри агрегата; outbox публикует его при save. Время агрегат получает параметром, а не берёт из
Instant.now(). Команда возвращает минимум: id или статус. Ни полный read-DTO, ни сам агрегат наружу не отдаются. - Валидация: контракт (Jakarta) — на контроллере; бизнес-инварианты — в методах агрегата.
- Версия агрегата растёт на каждую команду и пишется условием
WHERE version = ?; проигравший получает исключение, и обработчик либо повторяет идемпотентную команду, либо отдаёт409;FOR UPDATEдля постоянных конфликтов. - Ключ идемпотентности вставляют первым действием транзакции в таблицу с ответом; повтор с тем же хешем возвращает сохранённый ответ, с другим
422, в работе409; срок хранения не меньше суток, ключ привязан к актору. - Границы транзакции: внешние вызовы до или через outbox, откат только на непроверяемых исключениях, события после фиксации через
@TransactionalEventListener, вызов изнутри класса идёт мимо прокси. - Сага — последовательность локальных транзакций с компенсациями в обратном порядке; ей нужны идемпотентные шаги, своё состояние в таблице и сторож для застрявших, а упавшая компенсация требует повторов и ручного разбора.
- Когда решение зависит от чужого агрегата: передать готовое решение в команду, скопировать чужой факт событием, выполнить оптимистично с компенсацией или признать границы неверными — но не читать чужую таблицу и не звать чужой сервис внутри транзакции.
Что почитать дальше
- Query side в CQRS — как работает сторона чтения.
- Синхронизация через события — как outbox-событие из command-handler доходит до read-модели.
- Aggregate Root —
registerEvent, доменные методы, инварианты.