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

В 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-запрос контроллер команда record с данными диспетчер ищет обработчик обработчик транзакция открыта агрегат проверяет инварианты репозиторий save и версия таблица исходящих транзакция закрыта

Путь команды от 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}.

результат команды идентификатор создали объект пустой результат изменили состояние read-DTO смешивает чтение агрегат обход правил

Четыре возможных результата команды: сверху то, что вернуть можно, снизу то, что ломает разделение чтения и записи.

Решение зависит от чужого агрегата

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

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, вызов изнутри класса идёт мимо прокси.
  • Сага — последовательность локальных транзакций с компенсациями в обратном порядке; ей нужны идемпотентные шаги, своё состояние в таблице и сторож для застрявших, а упавшая компенсация требует повторов и ручного разбора.
  • Когда решение зависит от чужого агрегата: передать готовое решение в команду, скопировать чужой факт событием, выполнить оптимистично с компенсацией или признать границы неверными — но не читать чужую таблицу и не звать чужой сервис внутри транзакции.

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