Event-Driven архитектура: что это и зачем она нужна

Event-Driven архитектура: что это и зачем она нужна

Большинство backend-систем начинается одинаково: HTTP-запрос приходит, handler вызывает service, service дёргает repository, ответ уходит клиенту. Это request-response модель, и она отлично работает, пока у вас один сервис с простой логикой.

Проблемы начинаются, когда после основного действия нужно сделать три побочных: отправить email, обновить счётчик аналитики, начислить бонус. Каждый вызов добавляет латентность, и если одна из систем лежит, весь запрос падает. Event-Driven архитектура (EDA) решает именно эту проблему: основное действие выполняется быстро, а побочные эффекты обрабатываются асинхронно через события.

Request-Response vs Event-Driven

Sync handler знает обо всех side-effects и блокируется; EDA: один publish, подписчики работают параллельно

В синхронной модели сервис A напрямую вызывает сервис B и ждёт ответа. Каждый новый побочный эффект добавляет зависимость и время ожидания:

// Синхронный подход: handler знает обо всех побочных эффектах
func (h *LessonHandler) Complete(w http.ResponseWriter, r *http.Request) {
    lessonID := extractLessonID(r)
    userID := extractUserID(r)

    // Основное действие
    err := h.progressService.CompleteLessonForUser(r.Context(), userID, lessonID)
    if err != nil {
        http.Error(w, "failed to complete lesson", http.StatusInternalServerError)
        return
    }

    // Побочные эффекты - каждый добавляет латентность и точку отказа
    _ = h.emailService.SendLessonCompletedEmail(r.Context(), userID, lessonID)
    _ = h.analyticsService.TrackLessonCompletion(r.Context(), userID, lessonID)
    _ = h.bonusService.AwardPoints(r.Context(), userID, 10)

    w.WriteHeader(http.StatusOK)
}
<?php
// Синхронный подход: controller знает обо всех побочных эффектах
declare(strict_types=1);

#[Route('/lessons/{id}/complete', methods: ['POST'])]
public function complete(int $id, Request $request): Response
{
    $userId = (int) $request->attributes->get('userId');

    // Основное действие
    $this->progressService->completeLessonForUser($userId, $id);

    // Побочные эффекты - каждый добавляет латентность и точку отказа
    $this->emailService->sendLessonCompletedEmail($userId, $id);
    $this->analyticsService->trackLessonCompletion($userId, $id);
    $this->bonusService->awardPoints($userId, 10);

    return new Response(status: 200);
}

В event-driven модели handler выполняет основное действие и публикует событие. Подписчики реагируют независимо:

// Event-Driven подход: handler публикует событие, подписчики реагируют сами
func (h *LessonHandler) Complete(w http.ResponseWriter, r *http.Request) {
    lessonID := extractLessonID(r)
    userID := extractUserID(r)

    // Основное действие
    err := h.progressService.CompleteLessonForUser(r.Context(), userID, lessonID)
    if err != nil {
        http.Error(w, "failed to complete lesson", http.StatusInternalServerError)
        return
    }

    // Одна публикация - все подписчики обработают параллельно
    h.eventBus.Publish(events.LessonCompleted{
        UserID:   userID,
        LessonID: lessonID,
        At:       time.Now(),
    })

    w.WriteHeader(http.StatusOK)
}
<?php
// Event-Driven: controller публикует событие, подписчики реагируют сами
declare(strict_types=1);

use Symfony\Component\Messenger\MessageBusInterface;

#[Route('/lessons/{id}/complete', methods: ['POST'])]
public function complete(int $id, Request $request): Response
{
    $userId = (int) $request->attributes->get('userId');

    // Основное действие
    $this->progressService->completeLessonForUser($userId, $id);

    // Одна публикация - все подписчики обработают параллельно
    // (Symfony Messenger маршрутизирует событие подписчикам)
    $this->eventBus->dispatch(new LessonCompleted(
        userId: $userId,
        lessonId: $id,
        occurredAt: new \DateTimeImmutable(),
    ));

    return new Response(status: 200);
}

Handler больше не знает, кто слушает и что делает. Добавить новый побочный эффект (например, push-уведомление) - это добавить нового подписчика, не трогая handler.

Аналогия: ресторан

Представь ресторан. Официант (handler) принимает заказ от гостя и прикрепляет его на доску заказов (event bus). Повар, бариста и кондитер (subscribers) видят свои задачи и работают параллельно. Официант не стоит у каждого и не ждёт - он идёт к следующему столу.

Если ресторан нанимает нового повара на суши, официант ничего не меняет в своей работе - новый повар просто начинает читать заказы с доски. Это и есть слабая связность.

Ключевые понятия

Event (событие) - факт, который уже произошёл. Неизменяемый. Формулируется в прошедшем времени: LessonCompleted, UserRegistered, PaymentFailed.

Producer (издатель) - компонент, который публикует событие. Не знает, кто его прочитает.

Consumer (подписчик) - компонент, который реагирует на событие. Не знает, кто его отправил.

Channel / Broker - транспорт между producer и consumer. Может быть Go-каналом внутри процесса или внешним брокером (Kafka, RabbitMQ).

Простая event-система

Минимальный пример: publisher публикует событие, два subscriber-а реагируют на него.

package main

import (
    "fmt"
    "sync"
    "time"
)

// Event - базовое событие
type Event struct {
    Type      string
    Payload   map[string]any
    OccurredAt time.Time
}

// EventBus - простейшая шина событий на каналах
type EventBus struct {
    mu          sync.RWMutex
    subscribers map[string][]chan Event
}

func NewEventBus() *EventBus {
    return &EventBus{
        subscribers: make(map[string][]chan Event),
    }
}

// Subscribe регистрирует подписчика на определённый тип события
func (b *EventBus) Subscribe(eventType string) <-chan Event {
    b.mu.Lock()
    defer b.mu.Unlock()

    ch := make(chan Event, 16) // буферизованный канал
    b.subscribers[eventType] = append(b.subscribers[eventType], ch)
    return ch
}

// Publish отправляет событие всем подписчикам данного типа
func (b *EventBus) Publish(e Event) {
    b.mu.RLock()
    defer b.mu.RUnlock()

    for _, ch := range b.subscribers[e.Type] {
        // неблокирующая отправка: если подписчик не успевает - событие теряется
        select {
        case ch <- e:
        default:
            fmt.Printf("subscriber slow, event dropped: %s\n", e.Type)
        }
    }
}

func main() {
    bus := NewEventBus()

    // Подписчик 1: отправляет email
    emailCh := bus.Subscribe("lesson.completed")
    go func() {
        for e := range emailCh {
            fmt.Printf("[email] Урок %v пройден пользователем %v\n",
                e.Payload["lesson_id"], e.Payload["user_id"])
        }
    }()

    // Подписчик 2: обновляет аналитику
    analyticsCh := bus.Subscribe("lesson.completed")
    go func() {
        for e := range analyticsCh {
            fmt.Printf("[analytics] +1 completion для урока %v\n",
                e.Payload["lesson_id"])
        }
    }()

    // Публикация события
    bus.Publish(Event{
        Type:       "lesson.completed",
        Payload:    map[string]any{"user_id": 42, "lesson_id": 7},
        OccurredAt: time.Now(),
    })

    time.Sleep(100 * time.Millisecond) // даём горутинам отработать
}
<?php
// Простейшая in-process шина на массиве handler-ов
declare(strict_types=1);

final readonly class LessonCompleted
{
    public function __construct(
        public int $userId,
        public int $lessonId,
        public \DateTimeImmutable $occurredAt,
    ) {}
}

final class EventBus
{
    /** @var array<string, list<callable>> */
    private array $subscribers = [];

    public function subscribe(string $eventClass, callable $handler): void
    {
        $this->subscribers[$eventClass][] = $handler;
    }

    public function publish(object $event): void
    {
        foreach ($this->subscribers[$event::class] ?? [] as $handler) {
            $handler($event);
        }
    }
}

// Подписчики - просто invokable классы
final class SendEmailOnLessonCompleted
{
    public function __invoke(LessonCompleted $event): void
    {
        echo "[email] урок {$event->lessonId} пройден пользователем {$event->userId}\n";
    }
}

final class TrackAnalyticsOnLessonCompleted
{
    public function __invoke(LessonCompleted $event): void
    {
        echo "[analytics] +1 completion для урока {$event->lessonId}\n";
    }
}

$bus = new EventBus();
$bus->subscribe(LessonCompleted::class, new SendEmailOnLessonCompleted());
$bus->subscribe(LessonCompleted::class, new TrackAnalyticsOnLessonCompleted());

$bus->publish(new LessonCompleted(
    userId: 42,
    lessonId: 7,
    occurredAt: new \DateTimeImmutable(),
));

// В Symfony та же роль у Messenger:
// final class SendEmailOnLessonCompleted
// {
//     #[AsMessageHandler]
//     public function __invoke(LessonCompleted $event): void { ... }
// }

Поток событий в типичном приложении

Поток: HTTP → Use-Case → DB → Domain Event → Event Bus → независимые подписчики

Основное действие выполняется синхронно. Побочные эффекты - асинхронно. Если подписчик email упал, урок всё равно засчитан.

Когда использовать EDA

Хороший сценарий для событий - побочные эффекты, которые не влияют на основной результат:

  • Отправка email / push-уведомлений
  • Аналитика и метрики
  • Начисление бонусов / ачивок
  • Синхронизация данных между сервисами
  • Аудит-лог

Плохой сценарий - когда нужен немедленный, консистентный результат:

  • Простой CRUD без побочных эффектов (overhead не оправдан)
  • Действия с требованием строгой консистентности (перевод денег между счетами)
  • Запросы, где клиенту нужен результат побочного эффекта прямо сейчас
Event-Driven архитектура прекрасно работает внутри монолита. Если у вас один Go-сервис с тремя модулями (progress, notifications, analytics), EventBus на каналах развяжет их не хуже, чем Kafka между тремя микросервисами. Не нужно разбивать монолит, чтобы использовать события. С событиями сложнее отлаживать: нет единого call stack. Нужен correlation ID, чтобы отследить цепочку от HTTP-запроса до последнего подписчика. Без дисциплины именования и версионирования событий система быстро превращается в хаос.

Типичные ошибки

  • «Заэвентим всё» - каждое CRUD-действие → событие. Получаем хаос событий без бизнес-смысла, дебаг невозможен. Событие - про бизнес-факт (OrderPlaced, PaymentReceived), не про техническое действие (RowUpdated).

  • EDA вместо синхронного API «потому что модно» - пользователь жмёт «купить» и ждёт ответ. Event-driven здесь даст eventual consistency и сложность ради ничего. EDA уместен для побочных эффектов (отправь email, начисли бонусы) - не для основного потока.

  • Использовать события для строгой консистентности - «перевод денег между счетами через события» обычно даёт race conditions. Для атомарных операций - транзакция БД или Saga с компенсациями (урок 8), а не «просто событие».

  • Подписчик читает БД сразу после события и не находит данные - event опубликован внутри транзакции, ещё не закоммичено. Подписчик стартует, идёт в БД, видит «нет». Правильное место публикации - после commit, либо через Outbox (урок 7).

  • Один большой EventBus без bounded contexts - все события всех модулей в одной шине. Со временем каждый подписчик подписывается на каждое событие, EDA вырождается в RPC через посредника. Разделяй шины по domain area.

  • Долгая обработка в Publish-хендлере - Publish() блокирует publisher, пока все подписчики не отработают. Сделай Publish async (через канал или горутину пула) - иначе HTTP-запрос ждёт всех слушателей.

  • Python - Kafka с aiokafka: producer, consumer, partitions - конкретная реализация event-driven архитектуры с Kafka в Python

Мини-задание

  • Перечисли побочные эффекты в своём проекте (email, аналитика, нотификации, бонусы)
  • Определи, какие из них можно выполнять асинхронно, не блокируя основной запрос
  • Напиши простой EventBus на Go-каналах с методами Subscribe и Publish
  • Подключи два подписчика к одному типу события и убедись, что оба получают сообщение
  • Нарисуй поток данных своего приложения: где request-response, а где можно перейти на события

Зарегистрируйтесь бесплатно, чтобы пройти квиз, решить задание с автопроверкой и вести прогресс.