Domain Events: события домена и путь к Event-Driven

Domain Events: события домена и путь к Event-Driven

Пользователь прошёл урок. Что должно произойти дальше? Обновить прогресс, пересчитать процент завершения курса, отправить email "Поздравляем!", записать аналитику, проверить - может, пора выдать сертификат.

Если всю эту логику запихнуть в один метод CompleteLesson, получится бог-функция на 200 строк. Добавление каждой новой реакции потребует правки этого метода, а падение email-сервиса сломает основной сценарий.

Domain Events решают эту проблему: основной метод фиксирует факт ("урок пройден"), а все побочные реакции запускаются отдельно. Это базовый кирпичик event-driven архитектуры.

Событие - это факт в прошедшем времени

Domain Event - это неизменяемая запись о том, что уже произошло в домене. Не запрос, не команда, а свершившийся факт.

Command (Команда)Event (Событие)
ВремяБудущее, императивПрошедшее, свершившийся факт
ПримерCompleteLessonLessonCompleted
Может быть отклонена?ДаНет (факт уже случился)
Отправитель знает получателя?ДаНет (pub/sub)
ИзменяемостьМожет быть переформулированаИммутабельна

Команда говорит: "Сделай это". Событие говорит: "Это случилось".

События BackendStart

На платформе BackendStart можно выделить такие доменные события:

package domain

import "time"

// Пользователь завершил урок
type LessonCompleted struct {
    UserID     int64
    CourseSlug string
    LessonSlug string
    OccurredAt time.Time
}

// Пользователь прошёл квиз
type QuizPassed struct {
    UserID     int64
    LessonSlug string
    Score      int // процент правильных ответов
    OccurredAt time.Time
}

// Пользователь завершил весь курс
type CourseFinished struct {
    UserID     int64
    CourseSlug string
    OccurredAt time.Time
}

// Новый пользователь зарегистрировался
type UserRegistered struct {
    UserID    int64
    Email     string
    Provider  string // "github", "email"
    OccurredAt time.Time
}
<?php
declare(strict_types=1);

namespace App\Learning\Domain;

// Domain Event - final readonly: факт в прошедшем времени, иммутабельный.

// Пользователь завершил урок
final readonly class LessonCompleted
{
    public function __construct(
        public int $userId,
        public string $courseSlug,
        public string $lessonSlug,
        public \DateTimeImmutable $occurredAt,
    ) {}
}

// Пользователь прошёл квиз
final readonly class QuizPassed
{
    public function __construct(
        public int $userId,
        public string $lessonSlug,
        public int $score, // процент правильных ответов
        public \DateTimeImmutable $occurredAt,
    ) {}
}

// Пользователь завершил весь курс
final readonly class CourseFinished
{
    public function __construct(
        public int $userId,
        public string $courseSlug,
        public \DateTimeImmutable $occurredAt,
    ) {}
}

// Новый пользователь зарегистрировался
final readonly class UserRegistered
{
    public function __construct(
        public int $userId,
        public string $email,
        public string $provider, // 'github', 'email'
        public \DateTimeImmutable $occurredAt,
    ) {}
}

Обрати внимание на паттерн: имя в прошедшем времени, все поля неизменяемые, есть метка времени OccurredAt.

Структура события - это публичный контракт. Если контекст Analytics слушает LessonCompleted, любое удаление поля сломает подписчика. Добавлять поля безопасно, удалять - нет.

Publisher: интерфейс и реализация

Для публикации событий нужен простой интерфейс и in-memory реализация для старта:

package domain

// Интерфейс публикации событий
type EventPublisher interface {
    Publish(ctx context.Context, events ...interface{})
}
<?php
declare(strict_types=1);

namespace App\Learning\Domain;

// Интерфейс публикации событий
interface EventPublisher
{
    public function publish(object ...$events): void;
}
package infra

import (
    "context"
    "log/slog"
    "sync"
)

// EventHandler обрабатывает одно событие
type EventHandler interface {
    Handle(ctx context.Context, event interface{})
}

// InMemoryPublisher - синхронная реализация для монолита
type InMemoryPublisher struct {
    mu       sync.RWMutex
    handlers []EventHandler
}

func NewInMemoryPublisher() *InMemoryPublisher {
    return &InMemoryPublisher{}
}

func (p *InMemoryPublisher) Subscribe(h EventHandler) {
    p.mu.Lock()
    defer p.mu.Unlock()
    p.handlers = append(p.handlers, h)
}

func (p *InMemoryPublisher) Publish(ctx context.Context, events ...interface{}) {
    p.mu.RLock()
    defer p.mu.RUnlock()

    for _, event := range events {
        for _, h := range p.handlers {
            h.Handle(ctx, event)
        }
    }
}
<?php
declare(strict_types=1);

namespace App\Learning\Infrastructure;

use App\Learning\Domain\EventPublisher;

// EventHandler обрабатывает одно событие
interface EventHandler
{
    public function handle(object $event): void;
}

// InMemoryPublisher - синхронная реализация для монолита.
// В продакшене можно заменить на Symfony Messenger или RabbitMQ - интерфейс тот же.
final class InMemoryPublisher implements EventPublisher
{
    /** @var EventHandler[] */
    private array $handlers = [];

    public function subscribe(EventHandler $h): void
    {
        $this->handlers[] = $h;
    }

    public function publish(object ...$events): void
    {
        foreach ($events as $event) {
            foreach ($this->handlers as $h) {
                $h->handle($event);
            }
        }
    }
}

В монолите этого достаточно. Когда появится потребность в асинхронности, реализацию можно заменить на Kafka или NATS - интерфейс не изменится.

Агрегат собирает события

Агрегат не публикует события сразу. Он накапливает их в слайсе, а публикация происходит после успешного сохранения в базу:

package domain

type CourseProgress struct {
    ID         int64
    UserID     int64
    CourseSlug CourseSlug
    Lessons    []LessonProgress
    Percent    int

    events []interface{} // накопленные события
}

// CollectEvent добавляет событие в очередь
func (p *CourseProgress) CollectEvent(e interface{}) {
    p.events = append(p.events, e)
}

// FlushEvents возвращает накопленные события и очищает очередь
func (p *CourseProgress) FlushEvents() []interface{} {
    events := p.events
    p.events = nil
    return events
}

// CompleteLesson - бизнес-операция, которая генерирует событие
func (p *CourseProgress) CompleteLesson(lessonSlug string) error {
    for i, lp := range p.Lessons {
        if lp.Slug == lessonSlug {
            if lp.Completed {
                return nil // идемпотентность
            }
            p.Lessons[i].Completed = true
            p.recalcPercent()

            p.CollectEvent(LessonCompleted{
                UserID:     p.UserID,
                CourseSlug: string(p.CourseSlug),
                LessonSlug: lessonSlug,
                OccurredAt: time.Now(),
            })

            if p.Percent == 100 {
                p.CollectEvent(CourseFinished{
                    UserID:     p.UserID,
                    CourseSlug: string(p.CourseSlug),
                    OccurredAt: time.Now(),
                })
            }
            return nil
        }
    }
    return fmt.Errorf("lesson %s not found in course %s", lessonSlug, p.CourseSlug)
}
<?php
declare(strict_types=1);

namespace App\Learning\Domain;

// Aggregate Root - final, но НЕ readonly: меняются Lessons, Percent, events.
final class CourseProgress
{
    /**
     * @param LessonProgress[] $lessons
     * @param object[]         $events
     */
    public function __construct(
        public readonly int $id,
        public readonly int $userId,
        public readonly CourseSlug $courseSlug,
        private array $lessons,
        private int $percent,
        private array $events = [],
    ) {}

    // recordEvent добавляет событие в очередь
    private function recordEvent(object $e): void
    {
        $this->events[] = $e;
    }

    // pullEvents возвращает накопленные события и очищает очередь
    /** @return object[] */
    public function pullEvents(): array
    {
        $events = $this->events;
        $this->events = [];
        return $events;
    }

    // completeLesson - бизнес-операция, которая генерирует событие
    public function completeLesson(string $lessonSlug): void
    {
        foreach ($this->lessons as $i => $lp) {
            if ($lp->slug === $lessonSlug) {
                if ($lp->completed) {
                    return; // идемпотентность
                }
                $this->lessons[$i] = $lp->withCompleted(true);
                $this->recalcPercent();

                $this->recordEvent(new LessonCompleted(
                    userId: $this->userId,
                    courseSlug: $this->courseSlug->value(),
                    lessonSlug: $lessonSlug,
                    occurredAt: new \DateTimeImmutable(),
                ));

                if ($this->percent === 100) {
                    $this->recordEvent(new CourseFinished(
                        userId: $this->userId,
                        courseSlug: $this->courseSlug->value(),
                        occurredAt: new \DateTimeImmutable(),
                    ));
                }
                return;
            }
        }
        throw new \DomainException(sprintf(
            'lesson %s not found in course %s',
            $lessonSlug,
            $this->courseSlug->value(),
        ));
    }
}

Порядок в use-case: (1) загрузить агрегат, (2) вызвать бизнес-метод, (3) сохранить, (4) опубликовать события:

Aggregate накапливает события через CollectEvent, use case вызывает Save и затем FlushEvents в event bus, обработчики SendEmail/Stats/Cert работают независимо

func (uc *CompleteLessonUseCase) Execute(ctx context.Context, cmd CompleteLesson) error {
    progress, err := uc.progressRepo.FindByUserAndCourse(ctx, cmd.UserID, cmd.CourseSlug)
    if err != nil {
        return err
    }

    if err := progress.CompleteLesson(cmd.LessonSlug); err != nil {
        return err
    }

    if err := uc.progressRepo.Save(ctx, progress); err != nil {
        return err
    }

    // Публикуем ПОСЛЕ успешного сохранения
    uc.publisher.Publish(ctx, progress.FlushEvents()...)
    return nil
}
final class CompleteLessonUseCase
{
    public function __construct(
        private readonly ProgressRepository $progressRepo,
        private readonly EventPublisher $publisher,
    ) {}

    public function execute(CompleteLessonCommand $cmd): void
    {
        $progress = $this->progressRepo->findByUserAndCourse(
            $cmd->userId,
            $cmd->courseSlug,
        );

        $progress->completeLesson($cmd->lessonSlug);

        $this->progressRepo->save($progress);

        // Публикуем ПОСЛЕ успешного сохранения
        $this->publisher->publish(...$progress->pullEvents());
    }
}

Обработчики событий

Каждый обработчик выполняет одну побочную реакцию:

// Отправка email при завершении курса
type SendCourseFinishedEmail struct {
    mailer EmailSender
}

func (h *SendCourseFinishedEmail) Handle(ctx context.Context, event interface{}) {
    e, ok := event.(domain.CourseFinished)
    if !ok {
        return
    }
    _ = h.mailer.Send(ctx, e.UserID, "course-finished", map[string]string{
        "course": e.CourseSlug,
    })
}

// Обновление статистики при прохождении урока
type UpdateProgressStats struct {
    statsRepo StatsRepository
}

func (h *UpdateProgressStats) Handle(ctx context.Context, event interface{}) {
    e, ok := event.(domain.LessonCompleted)
    if !ok {
        return
    }
    _ = h.statsRepo.IncrementLessonsCompleted(ctx, e.UserID)
}
<?php
declare(strict_types=1);

namespace App\Learning\Application\Handler;

use App\Learning\Domain\CourseFinished;
use App\Learning\Domain\EmailSender;
use App\Learning\Domain\LessonCompleted;
use App\Learning\Domain\StatsRepository;
use App\Learning\Infrastructure\EventHandler;

// Отправка email при завершении курса
final class SendCourseFinishedEmail implements EventHandler
{
    public function __construct(
        private readonly EmailSender $mailer,
    ) {}

    public function handle(object $event): void
    {
        if (!$event instanceof CourseFinished) {
            return;
        }
        $this->mailer->send(
            userId: $event->userId,
            template: 'course-finished',
            params: ['course' => $event->courseSlug],
        );
    }
}

// Обновление статистики при прохождении урока
final class UpdateProgressStats implements EventHandler
{
    public function __construct(
        private readonly StatsRepository $statsRepo,
    ) {}

    public function handle(object $event): void
    {
        if (!$event instanceof LessonCompleted) {
            return;
        }
        $this->statsRepo->incrementLessonsCompleted($event->userId);
    }
}
Если email-сервис недоступен, урок всё равно должен считаться пройденным. Обработчик событий - это побочный эффект. В синхронной реализации оборачивай вызов в recover или логируй ошибку без возврата. В асинхронной - используй retry-очередь.

Eventual Consistency

В синхронном монолите события обрабатываются мгновенно. Но при переходе к очередям (Kafka, RabbitMQ) обработка становится eventually consistent: факт записан, а побочные эффекты придут позже.

Пример: пользователь завершил курс, но сертификат появится через 2-3 секунды, когда обработчик прочитает событие из очереди. Для пользователя это незаметно, но архитектурно - принципиальная разница.

Правило: если побочный эффект может быть задержан на секунды без ущерба для UX - используй события. Если нет (например, проверка прав доступа) - это не событие, а часть основного потока.

Event Store: аудит и воспроизведение

Event Store - хранилище всех доменных событий в хронологическом порядке. Это не замена базы данных, а дополнение: текущее состояние хранится в таблицах, а история изменений - в Event Store.

package domain

type StoredEvent struct {
    ID         int64
    EventType  string
    Payload    []byte // JSON
    OccurredAt time.Time
    AggregateID int64
}

type EventStore interface {
    Append(ctx context.Context, events ...StoredEvent) error
    LoadByAggregate(ctx context.Context, aggregateID int64) ([]StoredEvent, error)
}

Зачем это нужно:

  • Аудит: видно, кто и когда прошёл урок, даже если прогресс был сброшен
  • Отладка: можно воспроизвести последовательность действий пользователя
  • Аналитика: построить воронку прохождения курса по событиям

Для DDD Lite полноценный Event Sourcing избыточен. Достаточно простого лога событий в таблице domain_events - это даст аудит без сложности CQRS.

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

  • Определи 3 доменных события для своего проекта (имена в прошедшем времени)
  • Реализуй структуру одного события со всеми необходимыми полями
  • Напиши интерфейс EventPublisher и InMemoryPublisher
  • Добавь метод CollectEvent/FlushEvents в один из своих агрегатов
  • Реализуй обработчик для одного события (логирование или запись статистики)

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