Domain Events: события домена и путь к Event-Driven
Domain Events: события домена и путь к Event-Driven
Пользователь прошёл урок. Что должно произойти дальше? Обновить прогресс, пересчитать процент завершения курса, отправить email "Поздравляем!", записать аналитику, проверить - может, пора выдать сертификат.
Если всю эту логику запихнуть в один метод CompleteLesson, получится бог-функция на 200 строк. Добавление каждой новой реакции потребует правки этого метода, а падение email-сервиса сломает основной сценарий.
Domain Events решают эту проблему: основной метод фиксирует факт ("урок пройден"), а все побочные реакции запускаются отдельно. Это базовый кирпичик event-driven архитектуры.
Событие - это факт в прошедшем времени
Domain Event - это неизменяемая запись о том, что уже произошло в домене. Не запрос, не команда, а свершившийся факт.
| Command (Команда) | Event (Событие) | |
|---|---|---|
| Время | Будущее, императив | Прошедшее, свершившийся факт |
| Пример | CompleteLesson | LessonCompleted |
| Может быть отклонена? | Да | Нет (факт уже случился) |
| Отправитель знает получателя? | Да | Нет (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.
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) опубликовать события:
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);
}
}
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 в один из своих агрегатов
- Реализуй обработчик для одного события (логирование или запись статистики)