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-process | Integration-тесты, 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), бизнес-логика остаётся нетронутой.
Типичные ошибки
-
Брокер «потому что серьёзный проект» - 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)
- Напиши тест: опубликуй событие, проверь, что два подписчика оба получили его