Event Bus vs Message Broker: в чём разница

Когда ты решил использовать события, первый вопрос: как их доставлять? Есть два принципиально разных подхода. Event Bus - это in-process механизм: события передаются внутри одного приложения через каналы или callback-функции. Message Broker - это отдельный сервис (Kafka, RabbitMQ, NATS), который принимает, хранит и доставляет сообщения между разными процессами.

Выбор между ними - это не вопрос «что лучше», а вопрос масштаба задачи. Начинать стоит с самого простого решения, которое работает, и усложнять только когда появляются реальные ограничения.

Event Bus: события внутри процесса

Event Bus живёт в памяти приложения. Это обычная структура с подписчиками и методами Subscribe / Publish. Никакой сети, никаких зависимостей.

package eventbus

import (
    "log/slog"
    "sync"
)

// Handler - функция-обработчик события
type Handler func(event any)

// Bus - in-process шина событий
type Bus struct {
    mu       sync.RWMutex
    handlers map[string][]Handler
    logger   *slog.Logger
}

func New(logger *slog.Logger) *Bus {
    return &Bus{
        handlers: make(map[string][]Handler),
        logger:   logger,
    }
}

// Subscribe регистрирует обработчик для указанного типа события
func (b *Bus) Subscribe(eventType string, h Handler) {
    b.mu.Lock()
    defer b.mu.Unlock()
    b.handlers[eventType] = append(b.handlers[eventType], h)
    b.logger.Info("subscriber registered",
        slog.String("event_type", eventType),
    )
}

// Publish отправляет событие всем подписчикам в отдельных горутинах
func (b *Bus) Publish(eventType string, event any) {
    b.mu.RLock()
    handlers := b.handlers[eventType]
    b.mu.RUnlock()

    for i, h := range handlers {
        go func(idx int, handler Handler) {
            defer func() {
                if r := recover(); r != nil {
                    b.logger.Error("subscriber panicked",
                        slog.String("event_type", eventType),
                        slog.Int("subscriber_index", idx),
                        slog.Any("panic", r),
                    )
                }
            }()
            handler(event)
        }(i, h)
    }
}
<?php
// src/Infrastructure/EventBus/InProcessEventBus.php
declare(strict_types=1);

namespace App\Infrastructure\EventBus;

use Psr\Log\LoggerInterface;

interface EventBus
{
    public function subscribe(string $eventClass, callable $handler): void;

    public function publish(object $event): void;
}

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

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

    public function subscribe(string $eventClass, callable $handler): void
    {
        $this->handlers[$eventClass][] = $handler;
        $this->logger->info('subscriber registered', ['event_class' => $eventClass]);
    }

    public function publish(object $event): void
    {
        foreach ($this->handlers[$event::class] ?? [] as $idx => $handler) {
            try {
                $handler($event);
            } catch (\Throwable $e) {
                // Один сбойный подписчик не должен ронять остальных.
                $this->logger->error('subscriber failed', [
                    'event_class' => $event::class,
                    'subscriber_index' => $idx,
                    'err' => $e->getMessage(),
                ]);
            }
        }
    }
}

// В Symfony роль EventBus выполняет MessageBusInterface (Messenger).
// Подписчики помечаются #[AsMessageHandler] и автоматически связываются по типу.

Использование:

func main() {
    logger := slog.Default()
    bus := eventbus.New(logger)

    // Подписчик: отправка email
    bus.Subscribe("lesson.completed", func(event any) {
        e := event.(events.LessonCompleted)
        fmt.Printf("Sending email for user %d, lesson %d\n", e.UserID, e.LessonID)
    })

    // Подписчик: аналитика
    bus.Subscribe("lesson.completed", func(event any) {
        e := event.(events.LessonCompleted)
        fmt.Printf("Tracking completion for lesson %d\n", e.LessonID)
    })

    // Публикация
    bus.Publish("lesson.completed", events.LessonCompleted{
        UserID:   42,
        LessonID: 7,
    })
}
<?php
// public/index.php (или bin/console-команда)
declare(strict_types=1);

use App\Domain\Event\LessonCompleted;
use App\Infrastructure\EventBus\InProcessEventBus;
use Psr\Log\NullLogger;

require __DIR__ . '/../vendor/autoload.php';

$bus = new InProcessEventBus(new NullLogger());

// Подписчик: отправка email
$bus->subscribe(LessonCompleted::class, function (LessonCompleted $e): void {
    printf('Sending email for user %d, lesson %d', $e->userId, $e->lessonId);
});

// Подписчик: аналитика
$bus->subscribe(LessonCompleted::class, function (LessonCompleted $e): void {
    printf('Tracking completion for lesson %d', $e->lessonId);
});

// Публикация
$bus->publish(LessonCompleted::create(
    userId: 42,
    lessonId: 7,
    trackSlug: 'event-driven',
    timeSpentSec: 0,
    correlationId: 'demo',
));

Достоинства Bus: нулевая латентность, нет внешних зависимостей, легко тестировать. Недостатки: события теряются при перезапуске процесса, нельзя масштабировать подписчиков отдельно.

Message Broker: события между процессами

Message Broker - это отдельный сервис, который принимает сообщения от producer-ов и доставляет их consumer-ам. Основные игроки:

Apache Kafka - распределённый лог. Сообщения хранятся на диске, можно перечитать. Гарантирует порядок внутри partition. Подходит для высоких нагрузок (миллионы сообщений/сек).

RabbitMQ - классический брокер с очередями. Сообщения удаляются после обработки. Поддерживает сложный routing (exchanges, bindings). Подходит для задач с разными паттернами доставки.

NATS - минималистичный брокер. Очень быстрый, простой протокол. В базовой версии - at-most-once. NATS JetStream добавляет persistence и at-least-once.

Сравнение

ХарактеристикаEvent Bus (in-process)Message Broker (Kafka/RabbitMQ)
ЛатентностьНаносекундыМиллисекунды
DurabilityНет - всё в памятиДа - сообщения на диске
МасштабированиеТолько вертикальноеГоризонтальное (consumer groups)
Failure handlingПаника = потеря событийRetry, dead letter queue
Внешние зависимостиНетНужен отдельный сервис
ТестированиеUnit-тесты, всё in-processIntegration-тесты, testcontainers
OrderingГарантирован (один процесс)По partition (Kafka) / по очереди
ReplayНевозможенKafka: да, RabbitMQ: нет

Типизированный Event Bus для продакшна

В реальном проекте хочется типобезопасности. Вот Bus с generic-подписчиками:

package eventbus

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

// TypedHandler - типизированный обработчик
type TypedHandler[E any] func(ctx context.Context, event E) error

// Subscription хранит обработчик с типизацией через замыкание
type Subscription struct {
    handler func(ctx context.Context, event any) error
}

// TypedBus - шина с поддержкой контекста и ошибок
type TypedBus struct {
    mu     sync.RWMutex
    subs   map[string][]Subscription
    logger *slog.Logger
}

func NewTypedBus(logger *slog.Logger) *TypedBus {
    return &TypedBus{
        subs:   make(map[string][]Subscription),
        logger: logger,
    }
}

// Subscribe регистрирует типизированный обработчик
func Subscribe[E any](bus *TypedBus, eventType string, handler TypedHandler[E]) {
    bus.mu.Lock()
    defer bus.mu.Unlock()

    sub := Subscription{
        handler: func(ctx context.Context, event any) error {
            typed, ok := event.(E)
            if !ok {
                return fmt.Errorf("unexpected event type: %T", event)
            }
            return handler(ctx, typed)
        },
    }
    bus.subs[eventType] = append(bus.subs[eventType], sub)
}

// Publish вызывает всех подписчиков, логирует ошибки
func (b *TypedBus) Publish(ctx context.Context, eventType string, event any) {
    b.mu.RLock()
    subs := b.subs[eventType]
    b.mu.RUnlock()

    var wg sync.WaitGroup
    for _, sub := range subs {
        wg.Add(1)
        go func(s Subscription) {
            defer wg.Done()
            if err := s.handler(ctx, event); err != nil {
                b.logger.Error("subscriber failed",
                    slog.String("event_type", eventType),
                    slog.String("err", err.Error()),
                )
            }
        }(sub)
    }
    wg.Wait()
}
<?php
// src/Infrastructure/EventBus/TypedEventBus.php
declare(strict_types=1);

namespace App\Infrastructure\EventBus;

use Psr\Log\LoggerInterface;

// В PHP типобезопасность - через class-string и invokable handlers.
// Типизация подписчика на конкретный event-класс делается через signature `__invoke`.
final class TypedEventBus
{
    /** @var array<class-string, list<callable>> */
    private array $subs = [];

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

    /**
     * @template T of object
     * @param class-string<T> $eventClass
     * @param callable(T): void $handler
     */
    public function subscribe(string $eventClass, callable $handler): void
    {
        $this->subs[$eventClass][] = $handler;
    }

    public function publish(object $event): void
    {
        // PHP-FPM однопоточный per request - подписчики выполняются последовательно.
        // Для параллели используй Symfony Messenger с async-транспортом
        // (Redis/AMQP) и несколькими worker-процессами.
        foreach ($this->subs[$event::class] ?? [] as $handler) {
            try {
                $handler($event);
            } catch (\Throwable $e) {
                $this->logger->error('subscriber failed', [
                    'event_class' => $event::class,
                    'err' => $e->getMessage(),
                ]);
            }
        }
    }
}

Использование типизированной подписки:

bus := eventbus.NewTypedBus(slog.Default())

// Компилятор проверяет тип события
eventbus.Subscribe(bus, "lesson.completed",
    func(ctx context.Context, e events.LessonCompleted) error {
        // e уже типизирован - не нужен type assertion
        return sendEmail(ctx, e.UserID, e.LessonID)
    },
)
<?php
declare(strict_types=1);

use App\Domain\Event\LessonCompleted;
use App\Infrastructure\EventBus\TypedEventBus;

$bus = new TypedEventBus($logger);

// Сигнатура замыкания фиксирует тип события на этапе static analysis (Psalm/PHPStan)
$bus->subscribe(LessonCompleted::class, function (LessonCompleted $e) use ($emails): void {
    // $e уже типизирован - type assertion не нужен
    $emails->send($e->userId, $e->lessonId);
});

Стратегия миграции: Bus сегодня, Broker завтра

Правильный подход - определить интерфейс publisher/subscriber и начать с in-process реализации:

// Publisher - порт (интерфейс), не привязан к реализации
type Publisher interface {
    Publish(ctx context.Context, eventType string, event any) error
}

// Сегодня: InProcessPublisher (EventBus)
// Завтра: KafkaPublisher (без изменения use-case кода)
<?php
// src/Application/Port/EventPublisher.php
declare(strict_types=1);

namespace App\Application\Port;

// Publisher - порт. Use-case зависит от него, не от транспорта.
interface EventPublisher
{
    public function publish(object $event): void;
}

// Сегодня:   InProcessPublisher (Symfony Messenger sync-транспорт)
// Завтра:    AsyncPublisher (Messenger async-транспорт - Redis Streams / AMQP)
// Меняется одна строка в config/packages/messenger.yaml, бизнес-код не трогаем.

Use-case зависит от интерфейса Publisher, а не от конкретной реализации. Когда придёт время Kafka - меняется только DI-конфигурация (Wire), бизнес-логика остаётся нетронутой.

90% Go-монолитов прекрасно работают с in-process EventBus. Переходи на Kafka/RabbitMQ когда появится хотя бы одно из: (1) нужна гарантия доставки при перезапуске, (2) нужно масштабировать consumer-ы горизонтально, (3) появился второй сервис, которому нужны события первого. Kafka хранит сообщения, но это не замена PostgreSQL. Retention period по умолчанию - 7 дней. Не рассчитывай на Kafka как на permanent storage. Для event sourcing нужен отдельный event store.

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

  • Брокер «потому что серьёзный проект» - Kafka на старте маленького сервиса = месяцы пробоя operations (топологии, ZK/KRaft, ACL, мониторинг) при том же бизнес-результате, что EventBus на каналах. Брокер ставят, когда прижала потребность (durability, multi-consumer, бэкэп нагрузки).

  • In-process EventBus без panic recovery - panic в одном подписчике падает goroutine, дальше - фрагментарные обработки. Каждый subscriber обёрнут в defer recover() + лог. Без этого один баг в подписчике рушит соседей.

  • Использовать chan как «бесконечную очередь» - chan Event с буфером 1000 однажды заполнится; publisher заблокируется, HTTP-запросы зависнут. Чёткий контракт: либо ограниченный буфер + drop policy при overflow, либо очередь в БД.

  • Брокер выбран по «модности», не по семантике - Kafka хороша для upstream-логов и event sourcing (партиции, retention, replay), RabbitMQ - для рабочих очередей (per-message ack, routing, RPC-like). Выбор по фану выливается в боль через год.

  • Один топик/exchange на все события - events со всеми типами вперемешку. Consumer вынужден читать всё, фильтровать сам, плохо масштабируется. Разделяй по domain area или event type (orders.events, users.events).

  • EventBus опубликовал событие до commit транзакции - внутри Bus подписчик уже читает БД, видит «нет данных», падает. В отличие от брокера, EventBus не «подождёт коммит». Публикация - после успешного tx.Commit() (или через Outbox).

  • Тесты «всё в один EventBus» - параллельные тесты делят bus, подписчики одного теста ловят события другого. Создавай новый EventBus в t.Run / setup, не глобальный.

  • Go - Работа с базами данных - транзакции и Outbox pattern: связь между БД и брокером

  • Python - Kafka с aiokafka: broker в деталях - практика работы с Kafka как с message broker: topics, partitions, consumer groups

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

  • Реализуй EventBus с методами Subscribe(eventType, handler) и Publish(eventType, event)
  • Добавь recover в горутину подписчика, чтобы паника одного не убила весь процесс
  • Определи интерфейс Publisher, чтобы use-case не зависел от реализации (Bus или Broker)
  • Раздели свои события на две категории: внутренние (Bus) и межсервисные (Broker)
  • Напиши тест: опубликуй событие, проверь, что два подписчика оба получили его

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