Брокеры сообщений: AMQP/RabbitMQ vs Kafka vs Pulsar
В backend архитектуре часто нужно асинхронно передавать сообщения между сервисами. Например: пользователь зарегистрировался → отправить welcome email, обновить аналитику, проиндексировать в поиске. Сделать это в request-response модели медленно и хрупко. Решение - message broker. В этом уроке - зачем они нужны, какие бывают и как выбирать. Языко-независимый разбор той же темы есть в треке про брокеры сообщений.
Зачем нужен message broker
Сравним: пользователь регистрируется через POST /register - примерно такой эндпоинт мы писали в финальном проекте.
Без broker - синхронно:
@app.post("/register")
async def register(data: UserCreate):
user = await db.create_user(data)
await send_welcome_email(user) # 500ms
await update_analytics(user) # 300ms
await index_in_search(user) # 700ms
return user
# Общее время: 1500ms+, плюс если any service down - 500
С broker - асинхронно:
@app.post("/register")
async def register(data: UserCreate):
user = await db.create_user(data)
await broker.publish("user.registered", user_data) # 5ms
return user
# Время: 50ms, downstream services обрабатывают сами в своём ритме
Email-сервис, analytics, search consume событие из broker независимо. Если один down - сообщения накапливаются, обработаются когда поднимется. Это decoupling компонентов.
Основные концепции
- Producer (publisher) - кто шлёт сообщение
- Consumer (subscriber) - кто получает и обрабатывает
- Message - данные (обычно JSON, иногда protobuf)
- Topic / Queue / Stream - именованный канал для сообщений
- Broker - сервер хранения и доставки
Patterns:
- Pub/Sub - один producer, много subscribers (все получают копию)
- Work queue - один producer, один из consumers получает (load balancing)
- RPC - request-response через две очереди
RabbitMQ - классический AMQP broker
AMQP (Advanced Message Queuing Protocol) - стандарт messaging с гибкой routing моделью.
Producer → Exchange → (routing) → Queue → Consumer
- Exchange - принимает сообщения, решает в какие queues отправить
- Routing key - определяет маршрут
- Queue - буферизует до consumption
- Binding - связь exchange ↔ queue с правилом
Типы exchanges:
- direct - точное совпадение routing key
- topic - wildcard pattern (
order.created,order.*,#.error) - fanout - broadcast всем queues
- headers - routing по headers
Когда RabbitMQ:
- Work queues с complex routing
- RPC patterns
- Гарантированная доставка одному consumer
- Низкая latency на отдельных сообщениях
- Когда нужны acknowledgements и dead letter queues
Kafka - distributed event log
Совсем другая философия. Сообщения это events в append-only log:
Topic = log файл, разбитый на partitions
Producer → append → Partition
Consumer reads from offset
Особенности:
- Сообщения не удаляются при чтении - живут retention period (часы, дни, годы)
- Несколько consumer groups читают независимо со своими offsets
- Replay - можно перечитать историю с любого момента
- Очень высокий throughput (миллионы msg/sec)
- Partitions для параллельности
Когда Kafka:
- Event streaming, event sourcing
- Большие объёмы данных (terabytes/day)
- Многие consumer'ы одного потока (analytics + ML + audit)
- Replay capability важна
- Real-time data pipelines
Pulsar - современная альтернатива
Apache Pulsar объединяет лучшее из RabbitMQ и Kafka:
- Multi-tenancy из коробки
- Tiered storage (старые данные в S3)
- Geo-replication
- Schema registry
- Поддержка обеих моделей: queues и event streams
Меньше распространён в production, но активно растёт. Часто выбирается для new infrastructure.
Redis Streams
Redis с Streams (с 5.0) даёт легковесный messaging:
- Persistent streams в Redis
- Consumer groups как в Kafka
- При уже использовании Redis - меньше инфраструктуры
Хорошо для:
- Простых use cases
- Когда Redis уже в стеке
- Низкая операционная сложность
Limits: меньше throughput чем Kafka, нет sophisticated routing как у RabbitMQ.
Сравнение
| Аспект | RabbitMQ | Kafka | Pulsar | Redis Streams |
|---|---|---|---|---|
| Модель | Queues + routing | Event log | Both | Streams |
| Throughput | 50K msg/s | 1M+ msg/s | 1M+ | 100K msg/s |
| Latency | <10ms | 5-10ms | <10ms | <5ms |
| Persistence | Optional | Always | Always | Optional |
| Message replay | Нет | Да | Да | Да |
| Routing | Сложное | Простое | Сложное | Простое |
| Setup сложность | Средняя | Высокая | Средняя | Низкая |
| Use case | Tasks, RPC | Streams, events | Hybrid | Lightweight |
Нет универсального выбора - зависит от задачи.
Doctrine выбор
Простое правило:
- Task queues (process email, render PDF, send notification) → RabbitMQ или Celery
- Event streaming (audit logs, analytics, data pipeline) → Kafka
- Both? → Pulsar или Kafka с adapter
Не нужен broker:
- Маленький проект - может, не стоит overhead
- Synchronous OK - простые CRUD endpoints
- Memory queue достаточно - tasks внутри одного процесса (Celery с memory backend)
Гарантии доставки
| Гарантия | Описание | Использование |
|---|---|---|
| At-most-once | Может потеряться, не повторится | Метрики где OK потерять |
| At-least-once | Доставится, может повториться | Default для большинства |
| Exactly-once | Доставится ровно раз | Сложно, дорого |
Exactly-once часто иллюзия - стоит больше performance. На практике at-least-once + idempotent consumers (см. урок 59).
Дополнительные паттерны
Dead Letter Queue (DLQ) - сообщения которые не получилось обработать:
Original Queue → (failures) → DLQ → Manual inspection
Defines что делать с poison messages.
Outbox pattern - atomic save в БД + публикация:
Транзакция:
1. INSERT user
2. INSERT в outbox таблицу: {"event": "user.created", "data": {...}}
COMMIT
Background worker:
SELECT FROM outbox → publish to broker → DELETE from outbox
Гарантирует что либо и user сохранён и event опубликован, либо ничего. Альтернатива - dual write (две системы одновременно) - может разъехаться.
Сериализация
Брокеры передают bytes. Format:
- JSON - читаемо, schema implicit, медленнее
- Protobuf - бинарный, schema explicit, быстрее, нужны .proto файлы
- Avro - бинарный с schema registry, популярен в Kafka
- MessagePack - бинарный, проще proto
Для small scale - JSON OK. Для high-throughput - protobuf или Avro. Schema evolution критична для long-lived event streams.
Когда messaging не подходит
Не используй когда:
- Нужен синхронный ответ на запрос (use HTTP)
- Real-time с low latency требования (use direct connection)
- Маленький monolith без множества сервисов
- Команда не готова к operational complexity (брокеры требуют maintenance)
Messaging это добавочная инфраструктура - нужен monitoring, capacity planning, schema management. Не серебряная пуля.
Python библиотеки
| Broker | Python библиотека |
|---|---|
| RabbitMQ | aio-pika (async), pika (sync) |
| Kafka | aiokafka (async), kafka-python (sync), confluent-kafka (faster) |
| Pulsar | pulsar-client |
| Redis Streams | redis-py (async support) |
Также task queue frameworks на основе brokers:
- Celery - распространённый, поддерживает RabbitMQ, Redis
- ARQ - async-first, использует Redis
- Dramatiq - простой, поддерживает RabbitMQ, Redis
Task queues удобнее для tasks с retries, scheduling. Раздельные libraries для тонкого контроля над messaging.
Сравнение с Go и PHP
В Go: те же брокеры через свои клиенты. Самые популярные - segmentio/kafka-go (Kafka), streadway/amqp (RabbitMQ). Часто пишутся свои абстракции над низкоуровневыми клиентами.
В PHP: php-amqplib, php-rdkafka. Также enqueue PSR-style абстракция.
Брокеры независимы от языка - один RabbitMQ может обслуживать producers и consumers на Python, Go, PHP, JS одновременно. Это сила messaging - язык-agnostic.
Архитектурный пример
E-commerce платформа:
[User Service] →(user.registered)→ Kafka
↓
[Email Service] - send welcome
[Analytics] - track signup
[CRM] - create contact
[Search] - index user
[Order Service] →(order.placed)→ Kafka
↓
[Inventory] - decrement stock
[Notifications] - push to mobile
[Analytics] - track order
[Email] - confirmation
Каждый event обрабатывают многие consumers независимо. Order Service не знает про consumers - просто publishes. Декомпозиция через events.
Мини-задание
- Подними локально RabbitMQ через Docker:
docker run -d --name rabbit -p 5672:5672 -p 15672:15672 rabbitmq:3-management
# Management UI: http://localhost:15672 (guest/guest)
- Подними Kafka:
# docker-compose.yml для Kafka
version: "3"
services:
zookeeper:
image: confluentinc/cp-zookeeper:7.6.0
environment:
ZOOKEEPER_CLIENT_PORT: 2181
kafka:
image: confluentinc/cp-kafka:7.6.0
depends_on: [zookeeper]
ports:
- "9092:9092"
environment:
KAFKA_BROKER_ID: 1
KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
docker-compose up -d
- Подумай для своего проекта:
- Нужны ли тебе async messaging? Какие use cases?
- RabbitMQ или Kafka? Почему?
- Какие гарантии доставки приемлемы?
Что дальше
Освоили теорию messaging. В следующих двух уроках - конкретные реализации на Python: aio-pika для RabbitMQ и aiokafka для Kafka. После них финальный урок про consumer patterns.
Что дальше
Выбор сделан - дальше практика. В следующем уроке разберём RabbitMQ через aio-pika, затем Kafka через aiokafka, и завершим общими паттернами консумеров.