HardТеория5 min

CQRS и Event Sourcing

Разделение read/write моделей, event store, projections, snapshots. Когда использовать и когда НЕ использовать CQRS

CQRS -- Command Query Responsibility Segregation

CQRS -- паттерн, разделяющий модель на две части: команды (изменение состояния) и запросы (чтение состояния). Предложен Грегом Янгом (Greg Young) на основе CQS-принципа Бертрана Мейера.

Традиционный подход vs CQRS

    ТРАДИЦИОННЫЙ (CRUD)                    CQRS

    ┌──────────┐                  ┌──────────┐  ┌──────────┐
    │  Client  │                  │  Client  │  │  Client  │
    └────┬─────┘                  └────┬─────┘  └────┬─────┘
         │                             │              │
    ┌────▼─────┐                  ┌────▼─────┐  ┌────▼─────┐
    │  Single  │                  │ Command  │  │  Query   │
    │  Model   │                  │  Model   │  │  Model   │
    │(read+    │                  │(write)   │  │(read)    │
    │ write)   │                  └────┬─────┘  └────┬─────┘
    └────┬─────┘                       │              │
         │                        ┌────▼─────┐  ┌────▼─────┐
    ┌────▼─────┐                  │  Write   │  │  Read    │
    │ Database │                  │  Store   │  │  Store   │
    └──────────┘                  └──────────┘  └──────────┘

Зачем разделять

Проблема единой модели Решение CQRS
Чтение и запись требуют разных оптимизаций Отдельные модели для каждой задачи
Read-модель перегружена бизнес-правилами Query-модель -- простая проекция
Масштабирование всей модели целиком Независимое масштабирование read/write
Сложные JOIN для отображения Денормализованные read-модели

Архитектура CQRS

                         ┌──────────────┐
                         │    Client    │
                         └──────┬───────┘
                                │
                    ┌───────────┴───────────┐
                    │                       │
              ┌─────▼─────┐          ┌─────▼─────┐
              │  Command  │          │   Query   │
              │  Handler  │          │  Handler  │
              └─────┬─────┘          └─────┬─────┘
                    │                       │
              ┌─────▼─────┐          ┌─────▼─────┐
              │  Domain   │          │   Read    │
              │  Model    │          │   Model   │
              │(Aggregates│          │(Projec-   │
              │ Entities) │          │ tions)    │
              └─────┬─────┘          └─────┬─────┘
                    │                       │
              ┌─────▼─────┐          ┌─────▼─────┐
              │  Write    │   sync   │   Read    │
              │  Store    │─────────▶│   Store   │
              │(PostgreSQL│  async   │(Redis,    │
              │ normalized│          │ Elastic,  │
              │ tables)   │          │ denorm.)  │
              └───────────┘          └───────────┘

Command Side -- пример на PHP

<?php

declare(strict_types=1);

// Command: описывает НАМЕРЕНИЕ (intent)
final readonly class PlaceOrderCommand
{
    public function __construct(
        public string $customerId,
        public string $productId,
        public int $quantity,
        public string $shippingAddress,
    ) {}
}

// Command Handler: обрабатывает команду
final readonly class PlaceOrderHandler
{
    public function __construct(
        private OrderRepositoryInterface $orderRepository,
        private ProductRepositoryInterface $productRepository,
        private EventBusInterface $eventBus,
    ) {}

    public function handle(PlaceOrderCommand $command): string
    {
        $product = $this->productRepository->findOrFail($command->productId);

        // Business rules in Domain Model
        $order = Order::place(
            id: $this->orderRepository->nextIdentity(),
            customerId: $command->customerId,
            product: $product,
            quantity: $command->quantity,
            shippingAddress: $command->shippingAddress,
        );

        $this->orderRepository->save($order);

        // Publish domain events for read model sync
        foreach ($order->pullDomainEvents() as $event) {
            $this->eventBus->publish($event);
        }

        return (string) $order->getId();
    }
}

Query Side -- пример на PHP

<?php

declare(strict_types=1);

// Query: запрос на получение данных
final readonly class GetOrderDetailsQuery
{
    public function __construct(
        public string $orderId,
    ) {}
}

// Read Model: плоская денормализованная структура
final readonly class OrderDetailsView
{
    public function __construct(
        public string $orderId,
        public string $customerName,
        public string $productName,
        public int $quantity,
        public int $totalAmount,
        public string $status,
        public string $shippingAddress,
        public string $createdAt,
    ) {}
}

// Query Handler: читает из оптимизированного read store
final readonly class GetOrderDetailsHandler
{
    public function __construct(
        private \PDO $readDb, // Separate read database connection
    ) {}

    public function handle(GetOrderDetailsQuery $query): ?OrderDetailsView
    {
        // Direct query to denormalized read model -- no JOINs needed
        $stmt = $this->readDb->prepare(
            'SELECT * FROM order_details_view WHERE order_id = :id'
        );
        $stmt->execute(['id' => $query->orderId]);
        $row = $stmt->fetch(\PDO::FETCH_ASSOC);

        if (!$row) {
            return null;
        }

        return new OrderDetailsView(
            orderId: $row['order_id'],
            customerName: $row['customer_name'],
            productName: $row['product_name'],
            quantity: (int) $row['quantity'],
            totalAmount: (int) $row['total_amount'],
            status: $row['status'],
            shippingAddress: $row['shipping_address'],
            createdAt: $row['created_at'],
        );
    }
}

Event Sourcing

Event Sourcing -- паттерн, при котором состояние хранится как последовательность событий, а не как текущий снимок.

Традиционное хранение vs Event Sourcing

    ТРАДИЦИОННОЕ                    EVENT SOURCING

    ┌────────────────┐           ┌────────────────────┐
    │ orders table   │           │ event_store table   │
    ├────────────────┤           ├────────────────────┤
    │ id: 123        │           │ 1: OrderPlaced     │
    │ status: shipped│           │ 2: OrderConfirmed  │
    │ total: 5000    │           │ 3: ItemAdded       │
    │ items: 3       │           │ 4: DiscountApplied │
    │ updated: today │           │ 5: OrderShipped    │
    └────────────────┘           └────────────────────┘

    Только текущее                Полная история
    состояние                     всех изменений

Event Store -- структура

CREATE TABLE event_store (
    id            UUID PRIMARY KEY DEFAULT gen_random_uuid(),
    aggregate_id  UUID NOT NULL,
    aggregate_type VARCHAR(255) NOT NULL,
    event_type    VARCHAR(255) NOT NULL,
    event_data    JSONB NOT NULL,
    metadata      JSONB DEFAULT '{}',
    version       INTEGER NOT NULL,
    created_at    TIMESTAMPTZ NOT NULL DEFAULT NOW(),

    UNIQUE (aggregate_id, version)
);

CREATE INDEX idx_event_store_aggregate
    ON event_store (aggregate_id, version);

Агрегат с Event Sourcing на PHP

<?php

declare(strict_types=1);

// Domain Event
final readonly class OrderPlaced
{
    public function __construct(
        public string $orderId,
        public string $customerId,
        public string $productId,
        public int $quantity,
        public int $unitPrice,
        public \DateTimeImmutable $occurredAt,
    ) {}
}

final readonly class OrderShipped
{
    public function __construct(
        public string $orderId,
        public string $trackingNumber,
        public \DateTimeImmutable $occurredAt,
    ) {}
}

// Aggregate rebuilt from events
final class Order
{
    private string $status = '';
    private string $customerId = '';
    private int $totalAmount = 0;
    private int $version = 0;
    private array $uncommittedEvents = [];

    private function __construct(
        private readonly string $id,
    ) {}

    // Factory: create new order (generates event)
    public static function place(
        string $id,
        string $customerId,
        string $productId,
        int $quantity,
        int $unitPrice,
    ): self {
        $order = new self($id);

        $order->apply(new OrderPlaced(
            orderId: $id,
            customerId: $customerId,
            productId: $productId,
            quantity: $quantity,
            unitPrice: $unitPrice,
            occurredAt: new \DateTimeImmutable(),
        ));

        return $order;
    }

    // Rebuild from event history
    public static function fromHistory(string $id, array $events): self
    {
        $order = new self($id);

        foreach ($events as $event) {
            $order->applyEvent($event);
            $order->version++;
        }

        return $order;
    }

    public function ship(string $trackingNumber): void
    {
        if ($this->status !== 'confirmed') {
            throw new \DomainException('Can only ship confirmed orders');
        }

        $this->apply(new OrderShipped(
            orderId: $this->id,
            trackingNumber: $trackingNumber,
            occurredAt: new \DateTimeImmutable(),
        ));
    }

    private function apply(object $event): void
    {
        $this->applyEvent($event);
        $this->uncommittedEvents[] = $event;
    }

    // Event handlers that modify internal state
    private function applyEvent(object $event): void
    {
        match ($event::class) {
            OrderPlaced::class => $this->onOrderPlaced($event),
            OrderShipped::class => $this->onOrderShipped($event),
            default => throw new \RuntimeException('Unknown event: ' . $event::class),
        };
    }

    private function onOrderPlaced(OrderPlaced $event): void
    {
        $this->status = 'placed';
        $this->customerId = $event->customerId;
        $this->totalAmount = $event->quantity * $event->unitPrice;
    }

    private function onOrderShipped(OrderShipped $event): void
    {
        $this->status = 'shipped';
    }

    public function getUncommittedEvents(): array
    {
        return $this->uncommittedEvents;
    }

    public function getVersion(): int { return $this->version; }
}

Projections (проекции)

Проекция -- это read model, построенная из потока событий.

    Event Store                    Projections
    ┌──────────────┐
    │ OrderPlaced  │──────┐
    │ OrderConfirm │──┐   │
    │ ItemAdded    │──┤   ├───▶ ┌──────────────────┐
    │ OrderShipped │──┤   │     │ order_details_view│
    └──────────────┘  │   │     │ (denormalized)    │
                      │   │     └──────────────────┘
                      │   │
                      │   └───▶ ┌──────────────────┐
                      │         │ customer_orders   │
                      │         │ (per customer)    │
                      │         └──────────────────┘
                      │
                      └───────▶ ┌──────────────────┐
                                │ daily_revenue     │
                                │ (analytics)       │
                                └──────────────────┘
<?php

declare(strict_types=1);

// Projection: rebuilds read model from events
final class OrderDetailsProjection
{
    public function __construct(
        private readonly \PDO $readDb,
    ) {}

    public function onOrderPlaced(OrderPlaced $event): void
    {
        $this->readDb->prepare(
            'INSERT INTO order_details_view
             (order_id, customer_id, total_amount, status, created_at)
             VALUES (:oid, :cid, :amount, :status, :created)'
        )->execute([
            'oid' => $event->orderId,
            'cid' => $event->customerId,
            'amount' => $event->quantity * $event->unitPrice,
            'status' => 'placed',
            'created' => $event->occurredAt->format('c'),
        ]);
    }

    public function onOrderShipped(OrderShipped $event): void
    {
        $this->readDb->prepare(
            'UPDATE order_details_view
             SET status = :status, tracking_number = :tracking
             WHERE order_id = :oid'
        )->execute([
            'status' => 'shipped',
            'tracking' => $event->trackingNumber,
            'oid' => $event->orderId,
        ]);
    }
}

Snapshots (снимки)

Проблема Event Sourcing: восстановление агрегата из 10 000 событий -- медленно. Решение -- снимки.

    Без снимков:
    Event 1 → Event 2 → ... → Event 10000 → Текущее состояние
    (медленно: проигрываем все 10000 событий)

    Со снимками:
    Snapshot (v9500) → Event 9501 → ... → Event 10000 → Текущее состояние
    (быстро: загружаем снимок + 500 событий)
<?php

declare(strict_types=1);

final readonly class EventSourcedRepository
{
    private const SNAPSHOT_INTERVAL = 100;

    public function __construct(
        private EventStore $eventStore,
        private SnapshotStore $snapshotStore,
    ) {}

    public function load(string $aggregateId): Order
    {
        // Try to load from snapshot first
        $snapshot = $this->snapshotStore->find($aggregateId);

        if ($snapshot) {
            $events = $this->eventStore->loadAfterVersion(
                $aggregateId,
                $snapshot->version
            );
            return Order::fromSnapshot($snapshot, $events);
        }

        // No snapshot: rebuild from all events
        $events = $this->eventStore->loadAll($aggregateId);
        return Order::fromHistory($aggregateId, $events);
    }

    public function save(Order $order): void
    {
        $events = $order->getUncommittedEvents();
        $this->eventStore->append($order->getId(), $events, $order->getVersion());

        // Create snapshot every N events
        if ($order->getVersion() % self::SNAPSHOT_INTERVAL === 0) {
            $this->snapshotStore->save($order->toSnapshot());
        }
    }
}

Когда НЕ использовать CQRS и Event Sourcing

CQRS -- антипоказания

Ситуация Почему НЕ CQRS
Простой CRUD Overhead не оправдан, Layered достаточно
Соотношение read/write ~ 50/50 Нет выигрыша от разделения
Маленькая команда (1-3 чел.) Сложность поддержки двух моделей
Строгая согласованность Eventual consistency между read/write -- проблема
Простая доменная модель Не нужны сложные write-модели

Event Sourcing -- антипоказания

Ситуация Почему НЕ ES
Нужно удалять данные (GDPR) Удаление событий нарушает журнал
Простой домен Overhead хранения истории не оправдан
Нет требований аудита Главное преимущество ES не используется
Команда не знакома с ES Крутая кривая обучения
Высокие требования к задержке чтения Eventual consistency не подходит

Кто использует

Компания Что используют Контекст
Банки Event Sourcing Полная аудит-трейл транзакций
Event Store Ltd Event Sourcing EventStoreDB -- специализированная БД
Axon Framework CQRS + ES Java-фреймворк для CQRS
LMAX Exchange CQRS Финансовая биржа, миллионы ops/sec
Walmart CQRS Разделение каталога и корзины

Проверь себя

Как Event Sourcing решает проблему аудита?

Какая главная проблема Event Sourcing при восстановлении агрегата из большого количества событий?

Почему Event Sourcing создаёт проблемы с GDPR (право на удаление)?

Почему CQRS плохо подходит для простых CRUD-приложений?

Что означает CQRS?