Что такое 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, масштабирование) перевешивают сложность.