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 | Разделение каталога и корзины |