Publisher/Subscriber: подписки без сюрпризов
Publisher/Subscriber: подписки без сюрпризов
Pub/Sub - это паттерн, где Publisher не знает, кто получит его сообщение, а Subscriber не знает, кто его отправил. Они связаны только через тему (topic) и схему события. Это даёт максимальную развязку: можно добавлять и убирать подписчиков без изменения кода publisher-а.
Но развязка - это не бесплатный обед. Без дисциплины в именовании топиков, версионировании схем и обработке ошибок система быстро становится непредсказуемой. В этом уроке разберём, как организовать pub/sub так, чтобы подписки работали без сюрпризов.
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-ы не мигрируют.
Порядок обработки событий
В системе с одним подписчиком порядок прост: события обрабатываются в порядке публикации. С несколькими подписчиками и параллельной обработкой порядок ломается.
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 возвращает
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 - поле
email→emails: [], старые подписчики падают на парсинге. Каждое изменение схемы - новый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