Брокеры сообщений: 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.

Сравнение

АспектRabbitMQKafkaPulsarRedis Streams
МодельQueues + routingEvent logBothStreams
Throughput50K msg/s1M+ msg/s1M+100K msg/s
Latency<10ms5-10ms<10ms<5ms
PersistenceOptionalAlwaysAlwaysOptional
Message replayНетДаДаДа
RoutingСложноеПростоеСложноеПростое
Setup сложностьСредняяВысокаяСредняяНизкая
Use caseTasks, RPCStreams, eventsHybridLightweight

Нет универсального выбора - зависит от задачи.

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 библиотеки

BrokerPython библиотека
RabbitMQaio-pika (async), pika (sync)
Kafkaaiokafka (async), kafka-python (sync), confluent-kafka (faster)
Pulsarpulsar-client
Redis Streamsredis-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.

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

  1. Подними локально RabbitMQ через Docker:
docker run -d --name rabbit -p 5672:5672 -p 15672:15672 rabbitmq:3-management

# Management UI: http://localhost:15672 (guest/guest)
  1. Подними 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
  1. Подумай для своего проекта:
  • Нужны ли тебе async messaging? Какие use cases?
  • RabbitMQ или Kafka? Почему?
  • Какие гарантии доставки приемлемы?

Что дальше

Освоили теорию messaging. В следующих двух уроках - конкретные реализации на Python: aio-pika для RabbitMQ и aiokafka для Kafka. После них финальный урок про consumer patterns.

Что дальше

Выбор сделан - дальше практика. В следующем уроке разберём RabbitMQ через aio-pika, затем Kafka через aiokafka, и завершим общими паттернами консумеров.

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