Publisher/Subscriber: подписки без сюрпризов

Publisher/Subscriber: подписки без сюрпризов

Pub/Sub - это паттерн, где Publisher не знает, кто получит его сообщение, а Subscriber не знает, кто его отправил. Они связаны только через тему (topic) и схему события. Это даёт максимальную развязку: можно добавлять и убирать подписчиков без изменения кода publisher-а.

Но развязка - это не бесплатный обед. Без дисциплины в именовании топиков, версионировании схем и обработке ошибок система быстро становится непредсказуемой. В этом уроке разберём, как организовать pub/sub так, чтобы подписки работали без сюрпризов.

Publisher → topic → N subscribers; конвенция топика и Envelope

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

Publisher - это порт в терминах Hexagonal Architecture. Use-case зависит от интерфейса, а не от конкретного транспорта:

package events

import "context"

// Publisher - интерфейс для публикации событий
type Publisher interface {
    // Publish отправляет событие в указанный топик
    Publish(ctx context.Context, topic string, event Event) error
}
<?php
// src/Application/Port/EventPublisherPort.php
declare(strict_types=1);

namespace App\Application\Port;

use App\Domain\Event\DomainEvent;

interface EventPublisherPort
{
    public function publish(string $topic, DomainEvent $event): void;
}

// В Symfony альтернатива - MessageBusInterface из Messenger:
// $this->messageBus->dispatch($event);
// Маршрутизация по транспортам (async/sync) настраивается в messenger.yaml.

Реализация для in-process шины:

package eventbus

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

// InProcessPublisher - реализация Publisher через in-memory bus
type InProcessPublisher struct {
    mu       sync.RWMutex
    handlers map[string][]EventHandler
    logger   *slog.Logger
}

type EventHandler struct {
    Name    string                                    // имя подписчика для логов
    Handle  func(ctx context.Context, event any) error
}

func NewInProcessPublisher(logger *slog.Logger) *InProcessPublisher {
    return &InProcessPublisher{
        handlers: make(map[string][]EventHandler),
        logger:   logger,
    }
}

// Publish вызывает всех подписчиков топика
func (p *InProcessPublisher) Publish(ctx context.Context, topic string, event any) error {
    p.mu.RLock()
    handlers := make([]EventHandler, len(p.handlers[topic]))
    copy(handlers, p.handlers[topic])
    p.mu.RUnlock()

    for _, h := range handlers {
        if err := h.Handle(ctx, event); err != nil {
            p.logger.Error("subscriber error",
                slog.String("topic", topic),
                slog.String("subscriber", h.Name),
                slog.String("err", err.Error()),
            )
            // Продолжаем: ошибка одного подписчика не блокирует остальных
        }
    }
    return nil
}
<?php
// src/Infrastructure/EventBus/InProcessPublisher.php
declare(strict_types=1);

namespace App\Infrastructure\EventBus;

use App\Application\Port\EventPublisherPort;
use App\Domain\Event\DomainEvent;
use Psr\Log\LoggerInterface;

final class EventSubscription
{
    public function __construct(
        public readonly string $name,
        /** @var callable(DomainEvent): void */
        public readonly mixed $handle,
    ) {}
}

final class InProcessPublisher implements EventPublisherPort
{
    /** @var array<string, list<EventSubscription>> */
    private array $handlers = [];

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

    public function subscribe(string $topic, EventSubscription $sub): void
    {
        $this->handlers[$topic][] = $sub;
    }

    public function publish(string $topic, DomainEvent $event): void
    {
        foreach ($this->handlers[$topic] ?? [] as $sub) {
            try {
                ($sub->handle)($event);
            } catch (\Throwable $e) {
                // Продолжаем: ошибка одного подписчика не блокирует остальных
                $this->logger->error('subscriber error', [
                    'topic' => $topic,
                    'subscriber' => $sub->name,
                    'err' => $e->getMessage(),
                ]);
            }
        }
    }
}

Обрати внимание: ошибка одного подписчика логируется, но не останавливает обработку остальными. Publisher не должен падать из-за проблем subscriber-а.

Subscriber: типизированные обработчики

Subscriber получает событие и выполняет побочный эффект. Хороший subscriber - это чистая функция с одной ответственностью:

package subscribers

import (
    "context"
    "log/slog"
)

// EmailOnLessonCompleted отправляет email при завершении урока
type EmailOnLessonCompleted struct {
    emailSender EmailSender
    logger      *slog.Logger
}

func NewEmailOnLessonCompleted(sender EmailSender, logger *slog.Logger) *EmailOnLessonCompleted {
    return &EmailOnLessonCompleted{
        emailSender: sender,
        logger:      logger,
    }
}

func (s *EmailOnLessonCompleted) Handle(ctx context.Context, event any) error {
    e, ok := event.(events.LessonCompleted)
    if !ok {
        return fmt.Errorf("unexpected event type: %T", event)
    }

    s.logger.Info("sending lesson completed email",
        slog.Int64("user_id", e.UserID),
        slog.Int64("lesson_id", e.LessonID),
        slog.String("correlation_id", e.CorrelationID),
    )

    return s.emailSender.Send(ctx, e.UserID, "Урок пройден!", fmt.Sprintf(
        "Поздравляем! Вы завершили урок %d в треке %s.", e.LessonID, e.TrackSlug,
    ))
}
<?php
// src/Application/Subscriber/SendWelcomeEmailOnLessonCompleted.php
declare(strict_types=1);

namespace App\Application\Subscriber;

use App\Application\Port\EmailSenderPort;
use App\Domain\Event\LessonCompleted;
use Psr\Log\LoggerInterface;
use Symfony\Component\Messenger\Attribute\AsMessageHandler;

// final class - subscriber не должен расширяться/мокироваться через наследование.
// __invoke + #[AsMessageHandler] - идиоматичный Symfony Messenger.
#[AsMessageHandler]
final class SendWelcomeEmailOnLessonCompleted
{
    public function __construct(
        private readonly EmailSenderPort $emails,
        private readonly LoggerInterface $logger,
    ) {}

    public function __invoke(LessonCompleted $event): void
    {
        $this->logger->info('sending lesson completed email', [
            'user_id' => $event->userId,
            'lesson_id' => $event->lessonId,
            'correlation_id' => $event->correlationId,
        ]);

        $this->emails->send(
            userId: $event->userId,
            subject: 'Урок пройден!',
            body: sprintf(
                'Поздравляем! Вы завершили урок %d в треке %s.',
                $event->lessonId,
                $event->trackSlug,
            ),
        );
    }
}

Регистрация подписчиков при запуске приложения:

func setupSubscribers(bus *eventbus.InProcessPublisher, deps Dependencies) {
    // Email при завершении урока
    emailSub := subscribers.NewEmailOnLessonCompleted(deps.EmailSender, deps.Logger)
    bus.Subscribe("lesson.completed", eventbus.EventHandler{
        Name:   "email_on_lesson_completed",
        Handle: emailSub.Handle,
    })

    // Аналитика при завершении урока
    analyticsSub := subscribers.NewAnalyticsOnLessonCompleted(deps.Analytics, deps.Logger)
    bus.Subscribe("lesson.completed", eventbus.EventHandler{
        Name:   "analytics_on_lesson_completed",
        Handle: analyticsSub.Handle,
    })

    // Бонус при завершении курса
    bonusSub := subscribers.NewBonusOnCourseFinished(deps.BonusService, deps.Logger)
    bus.Subscribe("course.finished", eventbus.EventHandler{
        Name:   "bonus_on_course_finished",
        Handle: bonusSub.Handle,
    })
}
<?php
// В Symfony подписчики регистрируются автоматически через атрибут #[AsMessageHandler]
// и autoconfigure. Routing настраивается в config/packages/messenger.yaml.
//
// Для голого PHP без Symfony - явная регистрация:
declare(strict_types=1);

use App\Infrastructure\EventBus\EventSubscription;
use App\Infrastructure\EventBus\InProcessPublisher;

function setupSubscribers(InProcessPublisher $bus, Dependencies $deps): void
{
    // Email при завершении урока
    $emailSub = new SendWelcomeEmailOnLessonCompleted($deps->emails, $deps->logger);
    $bus->subscribe('lesson.completed', new EventSubscription(
        name: 'email_on_lesson_completed',
        handle: $emailSub(...),
    ));

    // Аналитика при завершении урока
    $analyticsSub = new AnalyticsOnLessonCompleted($deps->analytics, $deps->logger);
    $bus->subscribe('lesson.completed', new EventSubscription(
        name: 'analytics_on_lesson_completed',
        handle: $analyticsSub(...),
    ));

    // Бонус при завершении курса
    $bonusSub = new BonusOnCourseFinished($deps->bonus, $deps->logger);
    $bus->subscribe('course.finished', new EventSubscription(
        name: 'bonus_on_course_finished',
        handle: $bonusSub(...),
    ));
}

Топики и routing

Топик - это имя канала, по которому publisher и subscriber находят друг друга. Хорошая конвенция именования:

{domain}.{entity}.{action}

Примеры:
  progress.lesson.completed
  progress.quiz.passed
  auth.user.registered
  payment.subscription.renewed

Для in-process Bus топик - просто строка-ключ в map. Для Kafka - это реальный topic с partitions. Для RabbitMQ - routing key, привязанный к exchange.

Схема события: JSON с версионированием

Каждое событие - это контракт. Publisher и subscriber должны договориться о формате. Лучший способ - описать схему явно и добавить версию:

// Envelope - стандартная обёртка события
type Envelope struct {
    // Метаданные (одинаковые для всех событий)
    ID            string          `json:"id"`
    Type          string          `json:"type"`            // "progress.lesson.completed"
    Version       int             `json:"version"`         // 1, 2, 3...
    CorrelationID string          `json:"correlation_id"`
    OccurredAt    string          `json:"occurred_at"`     // RFC3339

    // Payload (специфичный для типа)
    Payload       json.RawMessage `json:"payload"`
}
<?php
// src/Infrastructure/Event/Envelope.php
declare(strict_types=1);

namespace App\Infrastructure\Event;

// В Symfony Messenger schema-versioning делается через Stamps, например
// SerializedMessageStamp + кастомный VersionStamp. Здесь - явная обёртка.
final readonly class Envelope
{
    public function __construct(
        // Метаданные (одинаковые для всех событий)
        public string $id,
        public string $type,            // 'progress.lesson.completed'
        public int $version,            // 1, 2, 3...
        public string $correlationId,
        public string $occurredAt,      // RFC3339

        // Payload (специфичный для типа) - JSON-строка
        public string $payload,
    ) {}
}

Пример JSON, который передаётся через broker:

{
    "id": "550e8400-e29b-41d4-a716-446655440000",
    "type": "progress.lesson.completed",
    "version": 2,
    "correlation_id": "req-abc-123",
    "occurred_at": "2026-05-03T14:30:00Z",
    "payload": {
        "user_id": 42,
        "lesson_id": 7,
        "track_slug": "event-driven",
        "time_spent_sec": 1200
    }
}

Эволюция схемы

Схемы событий неизбежно меняются. Правила безопасного изменения:

Backward compatible (безопасно):

  • Добавить новое поле с default-значением
  • Сделать обязательное поле опциональным
// Version 1
type LessonCompletedV1 struct {
    UserID   int64 `json:"user_id"`
    LessonID int64 `json:"lesson_id"`
}

// Version 2 - добавлено поле track_slug (старые consumer-ы его игнорируют)
type LessonCompletedV2 struct {
    UserID    int64  `json:"user_id"`
    LessonID  int64  `json:"lesson_id"`
    TrackSlug string `json:"track_slug"` // новое поле
}
<?php
declare(strict_types=1);

namespace App\Domain\Event;

// Version 1
final readonly class LessonCompletedV1
{
    public function __construct(
        public int $userId,
        public int $lessonId,
    ) {}
}

// Version 2 - добавлено поле trackSlug с default null
// Старые consumer-ы видят поле как unknown в JSON и игнорируют
final readonly class LessonCompletedV2
{
    public function __construct(
        public int $userId,
        public int $lessonId,
        public ?string $trackSlug = null,   // новое поле, опциональное
    ) {}
}

Breaking change (опасно):

  • Удалить поле
  • Переименовать поле
  • Изменить тип поля

При breaking change создавай новый тип события (progress.lesson.completed.v3) и поддерживай старый, пока все consumer-ы не мигрируют.

Добавь поле `version` в Envelope с самого начала, даже если кажется, что схема не изменится. Когда придёт время менять формат, ты скажешь себе спасибо. Consumer десериализует payload в зависимости от version: `switch envelope.Version { case 1: ... case 2: ... }`.

Порядок обработки событий

В системе с одним подписчиком порядок прост: события обрабатываются в порядке публикации. С несколькими подписчиками и параллельной обработкой порядок ломается.

Per-entity ordering - самая частая потребность. События одного пользователя обрабатываются в порядке. События разных пользователей - параллельно. В Kafka это достигается partition key = user_id:

// Kafka producer: partition key = user_id
// Все события одного пользователя попадают в одну partition
// и обрабатываются одним consumer-ом в порядке
func (p *KafkaPublisher) Publish(ctx context.Context, topic string, e events.Event) error {
    key := fmt.Sprintf("user:%d", e.UserID)
    msg := kafka.Message{
        Key:   []byte(key),
        Value: marshalEvent(e),
    }
    return p.writer.WriteMessages(ctx, msg)
}
<?php
// src/Infrastructure/Messenger/PartitionKeyMiddleware.php
declare(strict_types=1);

namespace App\Infrastructure\Messenger;

use App\Domain\Event\LessonCompleted;
use Symfony\Component\Messenger\Envelope;
use Symfony\Component\Messenger\Middleware\MiddlewareInterface;
use Symfony\Component\Messenger\Middleware\StackInterface;

// В Symfony Messenger partition key прокидывается через transport-specific stamps:
// - AMQP: AmqpStamp с routingKey
// - Redis Streams: RedisReceivedStamp
// Здесь - middleware, который добавляет ключ партиционирования по user_id.
final class PartitionKeyMiddleware implements MiddlewareInterface
{
    public function handle(Envelope $envelope, StackInterface $stack): Envelope
    {
        $message = $envelope->getMessage();

        if ($message instanceof LessonCompleted) {
            $envelope = $envelope->with(new PartitionKeyStamp('user:' . $message->userId));
        }

        return $stack->next()->handle($envelope, $stack);
    }
}

Global ordering (все события строго по порядку) - очень дорогая гарантия. Одна partition, один consumer, никакого параллелизма. Нужна крайне редко.

Обработка ошибок в subscriber-ах

Что делать, если subscriber упал при обработке события?

// RetryingHandler оборачивает handler с retry-логикой
type RetryingHandler struct {
    inner      EventHandler
    maxRetries int
    logger     *slog.Logger
}

func (h *RetryingHandler) Handle(ctx context.Context, event any) error {
    var lastErr error
    for attempt := 0; attempt <= h.maxRetries; attempt++ {
        if attempt > 0 {
            // Экспоненциальная задержка: 100ms, 200ms, 400ms...
            delay := time.Duration(1<<uint(attempt-1)) * 100 * time.Millisecond
            time.Sleep(delay)
            h.logger.Warn("retrying event handling",
                slog.String("subscriber", h.inner.Name),
                slog.Int("attempt", attempt),
            )
        }

        lastErr = h.inner.Handle(ctx, event)
        if lastErr == nil {
            return nil
        }
    }

    h.logger.Error("event handling failed after retries",
        slog.String("subscriber", h.inner.Name),
        slog.Int("max_retries", h.maxRetries),
        slog.String("err", lastErr.Error()),
    )

    // Отправляем в dead letter queue для ручного разбора
    return h.sendToDeadLetter(event, lastErr)
}
<?php
// src/Infrastructure/EventBus/RetryingHandler.php
declare(strict_types=1);

namespace App\Infrastructure\EventBus;

use App\Application\Port\DeadLetterQueue;
use Psr\Log\LoggerInterface;

// В Symfony Messenger retry настраивается в messenger.yaml через retry_strategy:
// max_retries, delay, multiplier - и failed transport для DLQ.
// Здесь - ручная обёртка для случая без Messenger.
final class RetryingHandler
{
    public function __construct(
        private readonly string $subscriberName,
        /** @var callable(object): void */
        private readonly mixed $inner,
        private readonly int $maxRetries,
        private readonly LoggerInterface $logger,
        private readonly DeadLetterQueue $dlq,
    ) {}

    public function __invoke(object $event): void
    {
        $lastException = null;

        for ($attempt = 0; $attempt <= $this->maxRetries; $attempt++) {
            if ($attempt > 0) {
                // Экспоненциальная задержка: 100ms, 200ms, 400ms...
                usleep((1 << ($attempt - 1)) * 100_000);
                $this->logger->warning('retrying event handling', [
                    'subscriber' => $this->subscriberName,
                    'attempt' => $attempt,
                ]);
            }

            try {
                ($this->inner)($event);
                return;
            } catch (\Throwable $e) {
                $lastException = $e;
            }
        }

        $this->logger->error('event handling failed after retries', [
            'subscriber' => $this->subscriberName,
            'max_retries' => $this->maxRetries,
            'err' => $lastException?->getMessage(),
        ]);

        // Отправляем в dead letter queue для ручного разбора
        $this->dlq->push($event, $lastException);
    }
}

Dead Letter Queue (DLQ) - очередь, куда попадают сообщения, которые не удалось обработать после всех retry. Инженер разбирает DLQ вручную или полуавтоматически. В Kafka это отдельный топик (progress.lesson.completed.dlq), в RabbitMQ - встроенный механизм.

Если subscriber получит одно и то же событие дважды (retry, дубликат от брокера), результат должен быть таким же, как при одной обработке. Например: «начислить 10 бонусов за урок 7» - проверяй, начислялись ли уже бонусы за этот урок, прежде чем начислять.

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

  • Subscriber возвращает nil при ошибке, чтобы «не ретраить» - событие потеряно молча, баг найдут через неделю. Возвращай ошибку, отправляй в DLQ, логируй с correlation_id - чтобы был след.
  • Бесконечный retry на «постоянную» ошибку - payload невалиден (UUID не парсится), retry будет падать вечно. Различай транзиентные ошибки (network, БД temporarily down) - retry со exp backoff, и permanent (валидация, schema mismatch) - сразу в DLQ.
  • Subscriber делает sync HTTP-вызов наружу без таймаута и retry-policy - внешний API лёг, подписчики висят, отставание растёт. Любой external call в subscriber-е - с deadline + circuit breaker.
  • Параллельная обработка событий одного aggregate-а - два subscriber-а одновременно обрабатывают OrderUpdated для order=42, второй пишет старое состояние поверх нового. Партиционируй по aggregate_id (Kafka partition key), либо последовательно в рамках aggregate.
  • Schema breaking change без version-bump - поле emailemails: [], старые подписчики падают на парсинге. Каждое изменение схемы - новый version, поддерживай N-1 версию параллельно минимум один деплой-цикл.
  • Обработка события > acknowledge timeout брокера - RabbitMQ heartbeat 60s, обработка 90s → broker считает consumer мёртвым, requeue, дубль. Либо короткая обработка, либо heartbeat увеличен, либо подтверждай рано + используй inbox.
  • Все subscribers получают всё - wildcard subscription на events.*. Под нагрузкой каждое событие копируется N раз. Подписывайся на конкретные topic-ы/event_type.

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

  • Определи интерфейс Publisher с методом Publish(ctx, topic, event) error
  • Напиши два subscriber-а для одного типа события (email + analytics)
  • Опиши JSON-схему одного события с полями id, type, version, correlation_id, occurred_at, payload
  • Реализуй RetryingHandler, который оборачивает обычный handler с retry-логикой
  • Придумай конвенцию именования топиков для своего проекта в формате domain.entity.action

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