Domain Event: «факт случился» - и это полезно

Domain Event - это запись о том, что уже произошло в системе (как в DDD). Не просьба, не команда, а свершившийся факт. Когда пользователь завершил урок, система не «просит завершить урок» - она фиксирует: «урок завершён». Эта разница между командой и событием - фундамент всей event-driven архитектуры (см. также Command).

Зачем вообще выделять события в отдельные структуры? Потому что событие - это контракт между частями системы. Use-case публикует событие, подписчики реагируют. Если событие хорошо спроектировано, подписчиков можно добавлять и убирать без изменения основной логики.

Именование: прошедшее время

Главное правило: событие всегда в прошедшем времени. Оно описывает факт, который уже случился:

Правильно (Event)Неправильно (Command)
LessonCompletedCompleteLesson
UserRegisteredRegisterUser
PaymentFailedFailPayment
QuizPassedPassQuiz
CourseFinishedFinishCourse
PasswordChangedChangePassword

Команда (CompleteLesson) может быть отклонена: урок уже завершён, пользователь заблокирован, данные невалидны. Событие (LessonCompleted) не может быть отклонено - оно уже произошло.

Структура события

Хорошее событие содержит достаточно информации, чтобы подписчик мог обработать его без дополнительных запросов в базу:

package events

import "time"

// Event - интерфейс, который реализуют все доменные события
type Event interface {
    // EventType возвращает уникальное имя типа события
    EventType() string
    // OccurredAt возвращает момент возникновения
    OccurredAt() time.Time
}

// BaseEvent - общие поля для всех событий
type BaseEvent struct {
    ID            string    `json:"id"`             // UUID события
    Type          string    `json:"type"`           // "lesson.completed"
    AggregateID   string    `json:"aggregate_id"`   // ID сущности, породившей событие
    CorrelationID string    `json:"correlation_id"` // связь с исходным запросом
    Timestamp     time.Time `json:"occurred_at"`    // когда произошло
}

func (e BaseEvent) EventType() string    { return e.Type }
func (e BaseEvent) OccurredAt() time.Time { return e.Timestamp }
<?php
// src/Domain/Event/DomainEvent.php
declare(strict_types=1);

namespace App\Domain\Event;

interface DomainEvent
{
    public function eventType(): string;

    public function occurredAt(): \DateTimeImmutable;
}

// src/Domain/Event/EventMetadata.php
// Общие метаданные для всех событий - immutable
final readonly class EventMetadata
{
    public function __construct(
        public string $id,                  // UUID события
        public string $type,                // "lesson.completed"
        public string $aggregateId,         // ID сущности, породившей событие
        public string $correlationId,       // связь с исходным запросом
        public \DateTimeImmutable $occurredAt,
    ) {}
}

CorrelationID - это идентификатор, который связывает событие с исходным HTTP-запросом. Когда подписчик обрабатывает событие и пишет лог, CorrelationID позволяет найти всю цепочку: от запроса пользователя до последнего побочного эффекта.

Конкретное событие с payload

Каждое событие расширяет базовые метаданные специфичными для домена полями:

package events

import (
    "time"

    "github.com/google/uuid"
)

// LessonCompleted - пользователь завершил урок
type LessonCompleted struct {
    BaseEvent
    UserID     int64  `json:"user_id"`
    LessonID   int64  `json:"lesson_id"`
    TrackSlug  string `json:"track_slug"`
    TimeSpent  int    `json:"time_spent_sec"` // сколько секунд ушло
}

// NewLessonCompleted создаёт событие с заполненными метаданными
func NewLessonCompleted(userID, lessonID int64, trackSlug string, timeSpent int, correlationID string) LessonCompleted {
    return LessonCompleted{
        BaseEvent: BaseEvent{
            ID:            uuid.New().String(),
            Type:          "lesson.completed",
            AggregateID:   fmt.Sprintf("user:%d", userID),
            CorrelationID: correlationID,
            Timestamp:     time.Now(),
        },
        UserID:    userID,
        LessonID:  lessonID,
        TrackSlug: trackSlug,
        TimeSpent: timeSpent,
    }
}
<?php
// src/Domain/Event/LessonCompleted.php
declare(strict_types=1);

namespace App\Domain\Event;

use Symfony\Component\Uid\Uuid;

// Событие - неизменяемая запись факта. final readonly гарантирует immutability.
final readonly class LessonCompleted implements DomainEvent
{
    public function __construct(
        public string $id,
        public string $aggregateId,
        public string $correlationId,
        public \DateTimeImmutable $occurredAt,
        public int $userId,
        public int $lessonId,
        public string $trackSlug,
        public int $timeSpentSec,
    ) {}

    public static function create(
        int $userId,
        int $lessonId,
        string $trackSlug,
        int $timeSpentSec,
        string $correlationId,
    ): self {
        return new self(
            id: Uuid::v4()->toRfc4122(),
            aggregateId: 'user:' . $userId,
            correlationId: $correlationId,
            occurredAt: new \DateTimeImmutable(),
            userId: $userId,
            lessonId: $lessonId,
            trackSlug: $trackSlug,
            timeSpentSec: $timeSpentSec,
        );
    }

    public function eventType(): string
    {
        return 'lesson.completed';
    }

    public function occurredAt(): \DateTimeImmutable
    {
        return $this->occurredAt;
    }
}

Обрати внимание: TimeSpent включён в событие, хотя подписчику «email» это не нужно. Зато подписчику «analytics» - очень нужно. Каждый подписчик берёт из события те поля, которые ему важны.

Event vs Command

Событие и команда - два разных паттерна, и путать их опасно:

// Command - приказ что-то сделать. Может быть отклонён.
type CompleteLesson struct {
    UserID   int64
    LessonID int64
}

// Event - факт, который уже произошёл. Не может быть отклонён.
type LessonCompleted struct {
    BaseEvent
    UserID   int64
    LessonID int64
}
<?php
declare(strict_types=1);

// Command - приказ что-то сделать. Может быть отклонён. Обрабатывает один handler.
final readonly class CompleteLesson
{
    public function __construct(
        public int $userId,
        public int $lessonId,
    ) {}
}

// Event - факт, который уже произошёл. Не может быть отклонён. Может слушать N подписчиков.
final readonly class LessonCompleted implements DomainEvent
{
    public function __construct(
        public int $userId,
        public int $lessonId,
        public \DateTimeImmutable $occurredAt,
    ) {}

    public function eventType(): string { return 'lesson.completed'; }
    public function occurredAt(): \DateTimeImmutable { return $this->occurredAt; }
}

Команда обрабатывается одним handler-ом: CompleteLessonHandler. Событие обрабатывается многими подписчиками: email, analytics, bonus, audit. Если команду обработали два handler-а - это баг. Если событие получили два подписчика - это нормальная работа.

Паттерн: сущность собирает события

В DDD-стиле доменная сущность собирает события во время выполнения бизнес-логики, а use-case диспатчит их после сохранения:

package domain

import "time"

// Progress - агрегат прогресса пользователя
type Progress struct {
    UserID      int64
    CompletedAt map[int64]time.Time // lessonID -> когда завершён
    events      []events.Event       // накопленные события
}

// CompleteLesson - бизнес-логика завершения урока
func (p *Progress) CompleteLesson(lessonID int64, trackSlug string, timeSpent int, corrID string) error {
    if _, ok := p.CompletedAt[lessonID]; ok {
        return ErrLessonAlreadyCompleted
    }

    p.CompletedAt[lessonID] = time.Now()

    // Сущность собирает событие, но НЕ публикует его
    p.events = append(p.events, events.NewLessonCompleted(
        p.UserID, lessonID, trackSlug, timeSpent, corrID,
    ))

    return nil
}

// PopEvents забирает накопленные события и очищает список
func (p *Progress) PopEvents() []events.Event {
    out := p.events
    p.events = nil
    return out
}
<?php
// src/Domain/Aggregate/Progress.php
declare(strict_types=1);

namespace App\Domain\Aggregate;

use App\Domain\Event\DomainEvent;
use App\Domain\Event\LessonCompleted;
use App\Domain\Exception\LessonAlreadyCompletedException;

// Aggregate root - не final, чтобы поддерживать наследование при необходимости.
class Progress
{
    /** @var array<int, \DateTimeImmutable> lessonID -> когда завершён */
    private array $completedAt = [];

    /** @var list<DomainEvent> накопленные события */
    private array $events = [];

    public function __construct(
        private readonly int $userId,
    ) {}

    public function completeLesson(int $lessonId, string $trackSlug, int $timeSpentSec, string $correlationId): void
    {
        if (isset($this->completedAt[$lessonId])) {
            throw new LessonAlreadyCompletedException($lessonId);
        }

        $this->completedAt[$lessonId] = new \DateTimeImmutable();

        // Агрегат собирает событие, но НЕ публикует его
        $this->events[] = LessonCompleted::create(
            userId: $this->userId,
            lessonId: $lessonId,
            trackSlug: $trackSlug,
            timeSpentSec: $timeSpentSec,
            correlationId: $correlationId,
        );
    }

    /** @return list<DomainEvent> */
    public function popEvents(): array
    {
        $out = $this->events;
        $this->events = [];
        return $out;
    }
}

Use-case вызывает PopEvents() после успешного сохранения в базу:

func (uc *CompleteLessonUseCase) Execute(ctx context.Context, cmd CompleteLesson) error {
    progress, err := uc.repo.GetByUserID(ctx, cmd.UserID)
    if err != nil {
        return fmt.Errorf("get progress: %w", err)
    }

    // Бизнес-логика - может вернуть ошибку
    if err := progress.CompleteLesson(cmd.LessonID, cmd.TrackSlug, cmd.TimeSpent, cmd.CorrelationID); err != nil {
        return err
    }

    // Сохранение - может вернуть ошибку
    if err := uc.repo.Save(ctx, progress); err != nil {
        return fmt.Errorf("save progress: %w", err)
    }

    // Публикация событий - только после успешного сохранения
    for _, e := range progress.PopEvents() {
        uc.bus.Publish(e)
    }

    return nil
}
<?php
// src/Application/UseCase/CompleteLessonUseCase.php
declare(strict_types=1);

namespace App\Application\UseCase;

use App\Application\Port\ProgressRepository;
use App\Domain\Command\CompleteLessonCommand;
use Symfony\Component\Messenger\MessageBusInterface;

final class CompleteLessonUseCase
{
    public function __construct(
        private readonly ProgressRepository $repo,
        private readonly MessageBusInterface $eventBus,
    ) {}

    public function __invoke(CompleteLessonCommand $cmd): void
    {
        $progress = $this->repo->getByUserId($cmd->userId);

        // Бизнес-логика - может бросить исключение
        $progress->completeLesson(
            lessonId: $cmd->lessonId,
            trackSlug: $cmd->trackSlug,
            timeSpentSec: $cmd->timeSpentSec,
            correlationId: $cmd->correlationId,
        );

        // Сохранение - может бросить исключение Doctrine
        $this->repo->save($progress);

        // Публикация событий - только после успешного сохранения
        foreach ($progress->popEvents() as $event) {
            $this->eventBus->dispatch($event);
        }
    }
}

Почему публикация после Save? Если опубликовать событие до сохранения, а Save упадёт, подписчики обработают событие для несуществующего факта. Это называется event leak.

Событие - это запись истории. Его нельзя редактировать или удалять. Если нужно «отменить» факт, публикуй компенсирующее событие: `LessonCompletionReverted`. Это как в бухгалтерии: запись не стирают, а делают сторнирующую проводку.

Сериализация событий

Для передачи через брокер или сохранения в event store события нужно сериализовать. JSON - самый распространённый формат:

package events

import "encoding/json"

// Envelope - обёртка для передачи через брокер
type Envelope struct {
    ID            string          `json:"id"`
    Type          string          `json:"type"`
    AggregateID   string          `json:"aggregate_id"`
    CorrelationID string          `json:"correlation_id"`
    OccurredAt    string          `json:"occurred_at"`
    Version       int             `json:"version"`       // версия схемы события
    Payload       json.RawMessage `json:"payload"`       // специфичные данные
}

// Wrap оборачивает событие в Envelope для передачи
func Wrap(e LessonCompleted) (Envelope, error) {
    payload, err := json.Marshal(struct {
        UserID    int64  `json:"user_id"`
        LessonID  int64  `json:"lesson_id"`
        TrackSlug string `json:"track_slug"`
        TimeSpent int    `json:"time_spent_sec"`
    }{
        UserID:    e.UserID,
        LessonID:  e.LessonID,
        TrackSlug: e.TrackSlug,
        TimeSpent: e.TimeSpent,
    })
    if err != nil {
        return Envelope{}, fmt.Errorf("marshal payload: %w", err)
    }

    return Envelope{
        ID:            e.ID,
        Type:          e.Type,
        AggregateID:   e.AggregateID,
        CorrelationID: e.CorrelationID,
        OccurredAt:    e.Timestamp.Format(time.RFC3339),
        Version:       1,
        Payload:       payload,
    }, nil
}
<?php
// src/Infrastructure/Event/Envelope.php
declare(strict_types=1);

namespace App\Infrastructure\Event;

use App\Domain\Event\LessonCompleted;

// Envelope - обёртка для передачи через брокер. В Symfony Messenger
// похожий механизм называется так же - `Symfony\Component\Messenger\Envelope`
// со stamps (CorrelationIdStamp, BusNameStamp, TransportMessageIdStamp).
final readonly class Envelope
{
    public function __construct(
        public string $id,
        public string $type,
        public string $aggregateId,
        public string $correlationId,
        public string $occurredAt,          // RFC3339
        public int $version,                // версия схемы события
        public string $payload,             // JSON-строка со специфичными данными
    ) {}

    public static function wrap(LessonCompleted $event): self
    {
        $payload = json_encode([
            'user_id' => $event->userId,
            'lesson_id' => $event->lessonId,
            'track_slug' => $event->trackSlug,
            'time_spent_sec' => $event->timeSpentSec,
        ], JSON_THROW_ON_ERROR);

        return new self(
            id: $event->id,
            type: 'lesson.completed',
            aggregateId: $event->aggregateId,
            correlationId: $event->correlationId,
            occurredAt: $event->occurredAt->format(\DateTimeInterface::RFC3339),
            version: 1,
            payload: $payload,
        );
    }
}

json.RawMessage для payload позволяет десериализовать метаданные без знания конкретного типа события. Подписчик смотрит на Type, определяет структуру payload и десериализует его в нужный тип.

Когда пользователь нажимает «Завершить урок», создаётся HTTP-запрос с уникальным ID. Этот ID прокидывается в событие как CorrelationID. Подписчик email, получив событие, пишет в лог тот же CorrelationID. Теперь в Kibana/Grafana одним запросом можно найти всю цепочку: HTTP → use-case → event → email subscriber.

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

  • Имя события в настоящем/повелительном времени - CreateOrder, SendEmail. Это уже не событие, а команда. События - факты прошлого: OrderCreated, EmailSent. Если имя звучит как приказ - это Command, обрабатывается одним получателем.
  • Анемичные события без payload - UserUpdated{ID: 42}. Подписчик вынужден идти в БД за данными. Race: к моменту запроса данные уже могли измениться. Клади в payload состояние на момент события - то, что нужно подписчикам для обработки.
  • Слишком жирное событие («god event») - весь объект User со всеми полями. Подписчик сильно связывается со схемой агрегата, любое изменение ломает контракт. Клади минимум, нужный для бизнес-смысла события + ссылку (ID) для подгрузки деталей.
  • Версионирование событий забыто - выкатили OrderPlaced без поля version/schema_version. Через полгода добавили обязательное поле - старые consumer-ы падают. Закладывай version с первого дня (см. урок 4).
  • OccurredAt ставит подписчик - событие гуляло в очереди час, метка времени - момент получения, а не возникновения. Метку ставит источник в момент создания события, не консьюмер.
  • Изменяемое событие - после публикации в payload событие правят. У разных подписчиков разная картина. Событие - immutable record: если факт неверен, выпускай корректирующее событие (OrderCancelled, PriceCorrected), не редактируй старое.
  • EventID не уникален - UUID v4 одинаковый у двух retry-publish из-за бага → подписчики с idempotency drop-ом не видят повторов. Генерируй ID до первой попытки публикации, переиспользуй при retry, не генерируй новый.

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

  • Определи интерфейс Event с методами EventType() и OccurredAt()
  • Создай BaseEvent struct и три конкретных события для своего домена (прошедшее время)
  • Реализуй паттерн «сущность собирает события» с методом PopEvents()
  • Напиши функцию сериализации события в JSON-формат Envelope
  • Сформулируй 5 событий домена и к каждому - соответствующую команду

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