Domain Event: «факт случился» - и это полезно
Domain Event - это запись о том, что уже произошло в системе (как в DDD). Не просьба, не команда, а свершившийся факт. Когда пользователь завершил урок, система не «просит завершить урок» - она фиксирует: «урок завершён». Эта разница между командой и событием - фундамент всей event-driven архитектуры (см. также Command).
Зачем вообще выделять события в отдельные структуры? Потому что событие - это контракт между частями системы. Use-case публикует событие, подписчики реагируют. Если событие хорошо спроектировано, подписчиков можно добавлять и убирать без изменения основной логики.
Именование: прошедшее время
Главное правило: событие всегда в прошедшем времени. Оно описывает факт, который уже случился:
| Правильно (Event) | Неправильно (Command) |
|---|---|
LessonCompleted | CompleteLesson |
UserRegistered | RegisterUser |
PaymentFailed | FailPayment |
QuizPassed | PassQuiz |
CourseFinished | FinishCourse |
PasswordChanged | ChangePassword |
Команда (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.
Сериализация событий
Для передачи через брокер или сохранения в 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 и десериализует его в нужный тип.
Типичные ошибки
- Имя события в настоящем/повелительном времени -
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() - Создай
BaseEventstruct и три конкретных события для своего домена (прошедшее время) - Реализуй паттерн «сущность собирает события» с методом
PopEvents() - Напиши функцию сериализации события в JSON-формат
Envelope - Сформулируй 5 событий домена и к каждому - соответствующую команду