Observer: события и подписчики (внутри приложения)

Когда одно событие должно запустить пять реакций, проще всего написать всё в одну функцию - и пожалеть через месяц. Observer аккуратно разводит источник и подписчиков.

Проблема

Пользователь прошёл урок. Нужно:

  1. Записать прогресс в БД
  2. Отправить email
  3. Обновить ачивки
  4. Отправить событие в аналитику

Если всё это засунуть в один use-case:

func (s *LessonService) Complete(userID, lessonID int64) error {
    s.repo.MarkCompleted(userID, lessonID)
    s.mailer.SendCongrats(userID)
    s.achievements.Check(userID)
    s.analytics.Track("lesson_completed", userID)
    // ...ещё 3 действия через месяц
    return nil
}
<?php
declare(strict_types=1);

final class LessonService
{
    public function __construct(
        private readonly LessonRepository $repo,
        private readonly Mailer $mailer,
        private readonly AchievementsService $achievements,
        private readonly AnalyticsTracker $analytics,
    ) {}

    public function complete(int $userId, int $lessonId): void
    {
        $this->repo->markCompleted($userId, $lessonId);
        $this->mailer->sendCongrats($userId);
        $this->achievements->check($userId);
        $this->analytics->track('lesson_completed', $userId);
        // ...ещё 3 действия через месяц
    }
}
class LessonService {
  #repo;
  #mailer;
  #achievements;
  #analytics;

  constructor({ repo, mailer, achievements, analytics }) {
    this.#repo = repo;
    this.#mailer = mailer;
    this.#achievements = achievements;
    this.#analytics = analytics;
  }

  async complete(userId, lessonId) {
    await this.#repo.markCompleted(userId, lessonId);
    await this.#mailer.sendCongrats(userId);
    await this.#achievements.check(userId);
    this.#analytics.track('lesson_completed', userId);
    // ...ещё 3 действия через месяц
  }
}

Проблемы: use-case знает обо всех зависимостях (нарушение SRP). Добавление нового действия - правка Complete() (нарушение OCP). Тестирование - мок 5 сервисов. Один упавший обработчик валит весь вызов.

Решение: Observer

Observer - подписка на события. Издатель (publisher) публикует событие. Подписчики (subscribers) реагируют на него. Издатель не знает, кто и сколько подписчиков.

sequenceDiagram
    autonumber
    participant P as Publisher<br/>(OrderService)
    participant B as EventBus
    participant S1 as Subscriber 1<br/>(Mailer)
    participant S2 as Subscriber 2<br/>(Analytics)
    participant S3 as Subscriber 3<br/>(Inventory)

    Note over P,S3: подписка происходит при старте приложения
    S1->>B: Subscribe("order.created")
    S2->>B: Subscribe("order.created")
    S3->>B: Subscribe("order.created")

    Note over P,S3: позже, в runtime
    P->>B: Publish("order.created", payload)
    B->>S1: notify(payload)
    B->>S2: notify(payload)
    B->>S3: notify(payload)
    Note right of B: Publisher не знает,<br/>кто и сколько слушает

Шина событий

type Event any

type Handler func(Event)

type EventBus struct {
    subs map[string][]Handler
    mu   sync.RWMutex
}

func NewEventBus() *EventBus {
    return &EventBus{subs: map[string][]Handler{}}
}

func (b *EventBus) Subscribe(event string, h Handler) {
    b.mu.Lock()
    defer b.mu.Unlock()
    b.subs[event] = append(b.subs[event], h)
}

func (b *EventBus) Publish(event string, data Event) {
    b.mu.RLock()
    defer b.mu.RUnlock()
    for _, h := range b.subs[event] {
        h(data)
    }
}

Теперь use-case публикует событие, а не вызывает всех руками:

type LessonCompleted struct {
    UserID   int64
    LessonID int64
}

func (s *LessonService) Complete(userID, lessonID int64) error {
    if err := s.repo.MarkCompleted(userID, lessonID); err != nil {
        return err
    }
    s.bus.Publish("lesson.completed", LessonCompleted{
        UserID: userID, LessonID: lessonID,
    })
    return nil
}

Подписчики регистрируются при старте приложения:

bus.Subscribe("lesson.completed", func(e Event) {
    evt := e.(LessonCompleted)
    mailer.SendCongrats(evt.UserID)
})

bus.Subscribe("lesson.completed", func(e Event) {
    evt := e.(LessonCompleted)
    achievements.Check(evt.UserID)
})
final class EventBus
{
    /** @var array<string, list<callable>> */
    private array $subs = [];

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

    public function publish(string $event, mixed $data): void
    {
        foreach ($this->subs[$event] ?? [] as $handler) {
            $handler($data);
        }
    }
}

Use-case публикует событие, а не вызывает всех руками:

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

final class LessonService
{
    public function __construct(
        private readonly LessonRepository $repo,
        private readonly EventBus $bus,
    ) {}

    public function complete(int $userId, int $lessonId): void
    {
        $this->repo->markCompleted($userId, $lessonId);
        $this->bus->publish('lesson.completed', new LessonCompleted($userId, $lessonId));
    }
}

Подписчики регистрируются при старте приложения:

$bus->subscribe('lesson.completed', function (LessonCompleted $e) use ($mailer) {
    $mailer->sendCongrats($e->userId);
});

$bus->subscribe('lesson.completed', function (LessonCompleted $e) use ($achievements) {
    $achievements->check($e->userId);
});
class EventBus {
  #subs = new Map();

  subscribe(event, handler) {
    const list = this.#subs.get(event) ?? [];
    list.push(handler);
    this.#subs.set(event, list);
  }

  publish(event, data) {
    for (const handler of this.#subs.get(event) ?? []) {
      handler(data);
    }
  }
}

Use-case публикует событие, а не вызывает всех руками:

class LessonService {
  #repo;
  #bus;
  constructor(repo, bus) {
    this.#repo = repo;
    this.#bus = bus;
  }

  async complete(userId, lessonId) {
    await this.#repo.markCompleted(userId, lessonId);
    this.#bus.publish('lesson.completed', { userId, lessonId });
  }
}

Подписчики регистрируются при старте приложения:

bus.subscribe('lesson.completed', ({ userId }) => mailer.sendCongrats(userId));
bus.subscribe('lesson.completed', ({ userId }) => achievements.check(userId));

В Node.js есть встроенный events.EventEmitter - почти готовый Observer (on/emit). Свой EventBus пишут, если нужен дополнительный контроль: фильтрация, middleware, асинхронный режим, перехват ошибок.

Use-case зависит только от EventBus, а не от 5 сервисов. Добавить нового подписчика - одна строка при регистрации, use-case не трогаем.

Observer vs Pub/Sub: в чём разница

Observer (GoF)Pub/Sub
Внутри одного процессаМежду процессами/сервисами
Синхронный (обычно)Асинхронный (брокер сообщений)
Прямая подписка на объектПодписка через канал/топик
Простая реализацияRabbitMQ, Kafka, NATS

Observer - это «маленький Pub/Sub» внутри приложения. Если обработчики тяжёлые (HTTP-вызовы, запись в БД), стоит перейти к полноценной очереди.

Observer с channels в Go

В Go события можно реализовать через каналы - это Go-специфичная конструкция, в PHP роль канала с буфером играет внешняя очередь (Redis/RabbitMQ через Symfony Messenger):

type LessonEvent struct {
    UserID   int64
    LessonID int64
}

func StartWorker(ch <-chan LessonEvent) {
    for evt := range ch {
        log.Printf("lesson %d completed by user %d", evt.LessonID, evt.UserID)
    }
}

// Публикация
ch := make(chan LessonEvent, 100)
go StartWorker(ch)

ch <- LessonEvent{UserID: 1, LessonID: 42}

Канал - встроенный примитив Go для передачи данных между горутинами. По сути это «Observer из коробки» с буфером и блокировкой.

EventBus (map + callbacks)Channels
Несколько подписчиков на событиеОдин читатель на канал (или fan-out)
Синхронный вызовМожно асинхронно (буферизация)
Проще для простых случаевЛучше для concurrency
Нет backpressureБуфер = встроенный backpressure

Обработка ошибок в Observer

Что делать, если один подписчик упал? Варианты:

// Вариант 1: логировать и продолжать (обычно лучший выбор)
func (b *EventBus) Publish(event string, data Event) {
    b.mu.RLock()
    defer b.mu.RUnlock()
    for _, h := range b.subs[event] {
        func() {
            defer func() {
                if r := recover(); r != nil {
                    log.Printf("handler panic: %v", r)
                }
            }()
            h(data)
        }()
    }
}

// Вариант 2: критические обработчики - в use-case, некритические - через bus
<?php
declare(strict_types=1);

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

    public function __construct(
        private readonly LoggerInterface $logger,
    ) {}

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

    // Вариант 1: логировать и продолжать (обычно лучший выбор)
    public function publish(string $event, mixed $data): void
    {
        foreach ($this->subs[$event] ?? [] as $handler) {
            try {
                $handler($data);
            } catch (Throwable $e) {
                $this->logger->error('handler failed', [
                    'event' => $event,
                    'err'   => $e->getMessage(),
                ]);
            }
        }
    }
}

// Вариант 2: критические обработчики - в use-case, некритические - через bus
class EventBus {
  #subs = new Map();
  #logger;
  constructor(logger) { this.#logger = logger ?? console; }

  subscribe(event, handler) {
    const list = this.#subs.get(event) ?? [];
    list.push(handler);
    this.#subs.set(event, list);
  }

  // Вариант 1: логировать и продолжать (обычно лучший выбор)
  publish(event, data) {
    for (const handler of this.#subs.get(event) ?? []) {
      try {
        handler(data);
      } catch (err) {
        this.#logger.error?.('handler failed', { event, err: err?.message });
      }
    }
  }
}

// Вариант 2: критические обработчики - в use-case, некритические - через bus
Observer внутри процесса обычно синхронный - подписчики выполняются в том же потоке. Если обработчик может «зависнуть» (HTTP-вызов, тяжёлая БД) - запускай его в горутине или используй очередь.

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

Один подписчик - если на событие реагирует одна функция, проще вызвать её напрямую. Observer оправдан при 2+ подписчиках.

Критическая последовательность - если порядок обработчиков важен (сначала записать в БД, потом отправить email), Observer плохо подходит - порядок подписчиков не гарантирован.

Отладка - событийная модель сложнее отлаживать. Стек вызовов не покажет, кто подписан. Используй хорошее логирование.

Связь с другими паттернами

Observer + Command - подписчик получает команду и выполняет её
Observer + Mediator - Mediator координирует объекты, Observer - уведомляет
Observer + Factory - фабрика создаёт подписчиков при регистрации

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

  • Реализуй EventBus и подпиши 2 обработчика на событие user.registered
  • Добавь обработку паники в Publish - чтобы один упавший подписчик не валил остальных
  • Попробуй вариант с channels: один канал, 3 горутины-обработчика
  • Подумай: какие события в твоём проекте стоит вынести из use-case в подписчики?

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