HardТеория7 min

Event-Driven архитектура

Event Sourcing, CQRS, Pub/Sub, очереди сообщений. PHP примеры с RabbitMQ и реализация Kafka consumer

Что такое Event-Driven

Event-Driven архитектура (EDA) -- это подход, при котором компоненты системы взаимодействуют через события, а не через прямые вызовы. Когда что-то происходит, генерируется событие, и заинтересованные компоненты реагируют.

Синхронный vs Event-Driven

Аспект Синхронный (REST) Event-Driven
Связанность Высокая (знает вызываемый сервис) Низкая (не знает потребителей)
Отказоустойчивость Каскадные сбои Изолированные сбои
Масштабирование Единообразное Независимое per consumer
Latency Суммируется по цепочке Асинхронное, не блокирует
Отладка Простая (стек вызовов) Сложная (распределенные события)

Основные паттерны

1. Pub/Sub (Publish/Subscribe)

Издатель публикует события, подписчики получают те, на которые подписаны.

<?php

declare(strict_types=1);

/**
 * Pub/Sub pattern with RabbitMQ
 *
 * Publisher doesn't know who will consume the message.
 * Multiple consumers can process the same event independently.
 */

// Domain Event
final readonly class OrderPlacedEvent
{
    public function __construct(
        public string $eventId,
        public string $orderId,
        public string $userId,
        public float $totalAmount,
        public array $items,
        public \DateTimeImmutable $occurredAt,
    ) {}

    public function toJson(): string
    {
        return json_encode([
            'event_id' => $this->eventId,
            'event_type' => 'order.placed',
            'order_id' => $this->orderId,
            'user_id' => $this->userId,
            'total_amount' => $this->totalAmount,
            'items' => $this->items,
            'occurred_at' => $this->occurredAt->format('c'),
        ]);
    }
}

// Publisher: Order Service
final class OrderEventPublisher
{
    public function __construct(
        private readonly \AMQPChannel $channel,
        private readonly string $exchangeName = 'orders',
    ) {}

    public function publishOrderPlaced(OrderPlacedEvent $event): void
    {
        $this->channel->basic_publish(
            new \AMQPMessage(
                $event->toJson(),
                [
                    'content_type' => 'application/json',
                    'delivery_mode' => 2, // Persistent
                    'message_id' => $event->eventId,
                    'timestamp' => time(),
                ],
            ),
            $this->exchangeName,
            'order.placed', // Routing key
        );
    }
}

// Consumer: Inventory Service (reduces stock)
final class InventoryEventConsumer
{
    public function __construct(
        private readonly InventoryService $inventory,
        private readonly LoggerInterface $logger,
    ) {}

    public function handleOrderPlaced(string $messageBody): void
    {
        $data = json_decode($messageBody, true);

        try {
            foreach ($data['items'] as $item) {
                $this->inventory->reserve(
                    $item['product_id'],
                    $item['quantity'],
                );
            }

            $this->logger->info('Inventory reserved for order', [
                'order_id' => $data['order_id'],
            ]);
        } catch (InsufficientStockException $e) {
            // Publish compensation event
            $this->logger->error('Insufficient stock', [
                'order_id' => $data['order_id'],
                'error' => $e->getMessage(),
            ]);
            throw $e; // Message will be nacked and retried
        }
    }
}

// Consumer: Notification Service (sends email)
final class NotificationEventConsumer
{
    public function __construct(
        private readonly NotificationService $notifier,
    ) {}

    public function handleOrderPlaced(string $messageBody): void
    {
        $data = json_decode($messageBody, true);

        $this->notifier->sendOrderConfirmation(
            $data['user_id'],
            $data['order_id'],
            $data['total_amount'],
        );
    }
}

2. Message Queue (Point-to-Point)

В отличие от Pub/Sub, каждое сообщение обрабатывается ровно одним потребителем.

<?php

declare(strict_types=1);

/**
 * Message Queue pattern: each message processed by ONE consumer
 *
 * Use for: task distribution, work queues, background jobs
 */
final class WorkQueueProducer
{
    public function __construct(
        private readonly \AMQPChannel $channel,
    ) {}

    /**
     * Enqueue a task for background processing.
     * Only one worker will pick up each task.
     */
    public function enqueue(string $queueName, array $task): void
    {
        $message = new \AMQPMessage(
            json_encode($task),
            [
                'content_type' => 'application/json',
                'delivery_mode' => 2, // Persistent
                'message_id' => bin2hex(random_bytes(16)),
            ],
        );

        $this->channel->basic_publish($message, '', $queueName);
    }
}

/**
 * Worker that processes tasks from a queue.
 * Includes retry logic and dead letter handling.
 */
final class WorkQueueConsumer
{
    public function __construct(
        private readonly \AMQPChannel $channel,
        private readonly LoggerInterface $logger,
        private readonly int $maxRetries = 3,
    ) {}

    public function consume(string $queueName, callable $handler): void
    {
        // Configure prefetch: only send 1 message at a time
        // to distribute work fairly among workers
        $this->channel->basic_qos(0, 1, false);

        $callback = function (\AMQPMessage $msg) use ($handler): void {
            $data = json_decode($msg->getBody(), true);
            $retryCount = ($msg->get('application_headers')['x-retry-count'] ?? 0);

            try {
                $handler($data);

                // Success: acknowledge the message
                $msg->getChannel()->basic_ack($msg->getDeliveryTag());
            } catch (\Throwable $e) {
                $this->logger->error('Task processing failed', [
                    'error' => $e->getMessage(),
                    'retry_count' => $retryCount,
                    'data' => $data,
                ]);

                if ($retryCount < $this->maxRetries) {
                    // Reject and re-queue with incremented retry count
                    $msg->getChannel()->basic_nack($msg->getDeliveryTag(), false, true);
                } else {
                    // Max retries exceeded: reject without re-queue
                    // Message goes to dead letter queue
                    $msg->getChannel()->basic_nack($msg->getDeliveryTag(), false, false);
                }
            }
        };

        $this->channel->basic_consume($queueName, '', false, false, false, false, $callback);

        while ($this->channel->is_consuming()) {
            $this->channel->wait();
        }
    }
}

Event Sourcing

Вместо хранения текущего состояния, храним все события, которые привели к текущему состоянию.

<?php

declare(strict_types=1);

/**
 * Event Sourcing: store events, derive state
 *
 * Traditional: UPDATE orders SET status = 'shipped'
 * Event Sourcing: APPEND event {type: 'OrderShipped', orderId: 123}
 *
 * Pros:
 * - Complete audit trail
 * - Can rebuild state at any point in time
 * - Natural for event-driven systems
 *
 * Cons:
 * - Complex querying (need projections)
 * - Storage grows over time
 * - Eventual consistency for read models
 */

// Domain events
interface DomainEvent
{
    public function getAggregateId(): string;
    public function getEventType(): string;
    public function getOccurredAt(): \DateTimeImmutable;
    public function toPayload(): array;
}

final readonly class OrderCreatedEvent implements DomainEvent
{
    public function __construct(
        private string $orderId,
        private string $userId,
        private array $items,
        private \DateTimeImmutable $occurredAt,
    ) {}

    public function getAggregateId(): string { return $this->orderId; }
    public function getEventType(): string { return 'order.created'; }
    public function getOccurredAt(): \DateTimeImmutable { return $this->occurredAt; }
    public function toPayload(): array
    {
        return ['user_id' => $this->userId, 'items' => $this->items];
    }
}

final readonly class OrderPaidEvent implements DomainEvent
{
    public function __construct(
        private string $orderId,
        private string $paymentId,
        private float $amount,
        private \DateTimeImmutable $occurredAt,
    ) {}

    public function getAggregateId(): string { return $this->orderId; }
    public function getEventType(): string { return 'order.paid'; }
    public function getOccurredAt(): \DateTimeImmutable { return $this->occurredAt; }
    public function toPayload(): array
    {
        return ['payment_id' => $this->paymentId, 'amount' => $this->amount];
    }
}

// Event Store
final class EventStore
{
    public function __construct(
        private readonly \PDO $db,
    ) {}

    public function append(DomainEvent $event): void
    {
        $this->db->prepare(
            'INSERT INTO event_store (aggregate_id, event_type, payload, occurred_at, version)
             VALUES (:aggregateId, :eventType, :payload, :occurredAt,
                     (SELECT COALESCE(MAX(version), 0) + 1
                      FROM event_store
                      WHERE aggregate_id = :aggregateId2))'
        )->execute([
            'aggregateId' => $event->getAggregateId(),
            'aggregateId2' => $event->getAggregateId(),
            'eventType' => $event->getEventType(),
            'payload' => json_encode($event->toPayload()),
            'occurredAt' => $event->getOccurredAt()->format('Y-m-d H:i:s.u'),
        ]);
    }

    /**
     * Load all events for an aggregate to rebuild its state.
     */
    public function loadEvents(string $aggregateId): array
    {
        $stmt = $this->db->prepare(
            'SELECT event_type, payload, occurred_at, version
             FROM event_store
             WHERE aggregate_id = :aggregateId
             ORDER BY version ASC'
        );
        $stmt->execute(['aggregateId' => $aggregateId]);

        return $stmt->fetchAll(\PDO::FETCH_ASSOC);
    }
}

// Aggregate that rebuilds state from events
final class Order
{
    private string $id;
    private string $status;
    private float $total = 0;
    private array $items = [];

    public static function fromEvents(array $events): self
    {
        $order = new self();

        foreach ($events as $event) {
            $order->apply($event['event_type'], json_decode($event['payload'], true));
        }

        return $order;
    }

    private function apply(string $eventType, array $payload): void
    {
        match ($eventType) {
            'order.created' => $this->applyCreated($payload),
            'order.paid' => $this->applyPaid($payload),
            'order.shipped' => $this->status = 'shipped',
            'order.cancelled' => $this->status = 'cancelled',
            default => null, // Unknown events are ignored
        };
    }

    private function applyCreated(array $payload): void
    {
        $this->status = 'pending';
        $this->items = $payload['items'];
    }

    private function applyPaid(array $payload): void
    {
        $this->status = 'paid';
        $this->total = $payload['amount'];
    }
}

CQRS (Command Query Responsibility Segregation)

Разделение модели на команды (изменение) и запросы (чтение).

<?php

declare(strict_types=1);

/**
 * CQRS: separate write and read models
 *
 * Write Model: processes commands, enforces business rules
 * Read Model: optimized for queries, denormalized
 *
 * Sync between models: via events (eventually consistent)
 */

// COMMAND side: validates and executes business logic
final class PlaceOrderCommandHandler
{
    public function __construct(
        private readonly EventStore $eventStore,
        private readonly InventoryChecker $inventory,
        private readonly EventDispatcher $dispatcher,
    ) {}

    public function handle(PlaceOrderCommand $command): string
    {
        // Validate business rules
        foreach ($command->items as $item) {
            if (!$this->inventory->isAvailable($item['productId'], $item['quantity'])) {
                throw new InsufficientStockException($item['productId']);
            }
        }

        $orderId = $this->generateOrderId();
        $event = new OrderCreatedEvent(
            $orderId,
            $command->userId,
            $command->items,
            new \DateTimeImmutable(),
        );

        // Store event (source of truth)
        $this->eventStore->append($event);

        // Dispatch for read model update and other consumers
        $this->dispatcher->dispatch($event);

        return $orderId;
    }

    private function generateOrderId(): string
    {
        return sprintf('ord_%s', bin2hex(random_bytes(12)));
    }
}

// QUERY side: optimized read model
final class OrderQueryService
{
    public function __construct(
        private readonly \PDO $readDb,
    ) {}

    /**
     * Read from denormalized view.
     * Updated asynchronously by event handlers.
     */
    public function getOrderSummary(string $orderId): ?array
    {
        $stmt = $this->readDb->prepare("
            SELECT
                o.id, o.status, o.total_amount,
                o.item_count, o.user_name, o.user_email,
                o.created_at, o.last_updated_at,
                o.shipping_address, o.tracking_number
            FROM order_summaries o
            WHERE o.id = :orderId
        ");
        $stmt->execute(['orderId' => $orderId]);

        return $stmt->fetch(\PDO::FETCH_ASSOC) ?: null;
    }

    public function getUserOrders(string $userId, int $limit = 20): array
    {
        $stmt = $this->readDb->prepare("
            SELECT id, status, total_amount, item_count, created_at
            FROM order_summaries
            WHERE user_id = :userId
            ORDER BY created_at DESC
            LIMIT :limit
        ");
        $stmt->execute(['userId' => $userId, 'limit' => $limit]);

        return $stmt->fetchAll(\PDO::FETCH_ASSOC);
    }
}

// Projection: updates read model when events occur
final class OrderSummaryProjection
{
    public function __construct(
        private readonly \PDO $readDb,
    ) {}

    public function onOrderCreated(array $event): void
    {
        $payload = $event['payload'];

        $this->readDb->prepare("
            INSERT INTO order_summaries
                (id, user_id, status, item_count, total_amount, created_at)
            VALUES
                (:id, :userId, 'pending', :itemCount, 0, :createdAt)
        ")->execute([
            'id' => $event['aggregate_id'],
            'userId' => $payload['user_id'],
            'itemCount' => count($payload['items']),
            'createdAt' => $event['occurred_at'],
        ]);
    }

    public function onOrderPaid(array $event): void
    {
        $this->readDb->prepare("
            UPDATE order_summaries
            SET status = 'paid',
                total_amount = :amount,
                last_updated_at = :updatedAt
            WHERE id = :id
        ")->execute([
            'id' => $event['aggregate_id'],
            'amount' => $event['payload']['amount'],
            'updatedAt' => $event['occurred_at'],
        ]);
    }
}

Когда использовать Event-Driven

Сценарий Подходит Не подходит
Микросервисы Да -- loose coupling
Simple CRUD Да -- overkill
Audit trail обязателен Да -- event sourcing
Низкая латентность Да -- async overhead
Независимые команды Да -- each team owns consumers
Строгая консистентность Да -- eventual consistency

Выводы

Event-Driven архитектура -- мощный инструмент для построения масштабируемых и отказоустойчивых систем. Но она добавляет сложность: eventual consistency, дебаг распределенных событий, идемпотентность потребителей. Используйте EDA когда преимущества (loose coupling, масштабирование) перевешивают сложность.