🐰 RabbitMQ

Полный справочник: AMQP-модель, обменники и маршрутизация, очереди (classic, quorum, stream), подтверждения, DLX и TTL, паттерны, rabbitmqctl, политики, кластер, клиенты и диагностика.

Шпаргалки · Инфраструктура · #rabbitmq #messaging #amqp #queues

Что это

RabbitMQ брокер сообщений (AMQP 0-9-1, также MQTT, STOMP, AMQP 1.0). Продюсер публикует в обменник, тот по правилам раскладывает сообщения по очередям, консьюмеры читают очереди. Подходит для очередей задач, асинхронных вызовов, маршрутизации и распределения нагрузки.

Модель

Термин Смысл
Producer отправитель
Exchange принимает сообщения и маршрутизирует по привязкам
Binding правило «обменник → очередь» с ключом или шаблоном
Queue буфер сообщений до обработки
Consumer читатель очереди
Connection / Channel TCP-соединение / лёгкий канал внутри него (создавайте один канал на поток)
Virtual host (vhost) изолированное пространство имён и прав
Routing key ключ маршрутизации в сообщении

Продюсер никогда не пишет в очередь напрямую: пустой обменник (default) направляет по имени очереди как ключу.

Типы обменников

Тип Маршрутизация Пример
direct точное совпадение ключа error, info
fanout во все привязанные очереди рассылка событий
topic шаблон: * одно слово, # ноль и больше logs.*.error, orders.#
headers по заголовкам реже используется

Типы очередей

Тип Особенности
Classic по умолчанию, быстрые, без репликации (классические зеркальные устарели)
Quorum реплицируемые (Raft), безопасность данных: рекомендуются для важных очередей; поддерживают DLX, x-delivery-limit
Stream журнал как в Kafka: хранится, читается повторно, большая пропускная способность

Параметры очереди: durable (переживает перезапуск), exclusive (для одного соединения), auto-delete. Для сообщений: delivery_mode=2 (persistent). Долговечность требует обоих: durable-очередь и persistent-сообщение.

Надёжность

  1. Подтверждения публикации (publisher confirms): брокер подтверждает приём.
  2. Подтверждения потребителя (ack): сообщение удаляется после basic_ack; nack / reject с requeue=true|false возвращает или отбрасывает. Автоподтверждение (auto_ack) может терять сообщения.
  3. Prefetch (basic_qos): сколько неподтверждённых сообщений держит потребитель (1 для равномерного распределения тяжёлых задач, 10–100 для быстрых).
  4. Обработка идемпотентна: возможны повторные доставки (redelivered=true).
  5. Dead Letter Exchange (DLX): куда попадают отклонённые, просроченные и переполнившие очередь сообщения.
  6. TTL: время жизни сообщения (x-message-ttl) или очереди (x-expires).
  7. Лимиты: x-max-length, x-max-length-bytes, x-overflow (drop-head, reject-publish).
# Политики: задаются на сервере, без пересоздания очередей
rabbitmqctl set_policy dlx-orders "^orders\." '{"dead-letter-exchange":"dlx","message-ttl":60000}' --apply-to queues

Паттерны

Паттерн Реализация
Рабочая очередь (work queue) одна очередь, несколько потребителей, prefetch=1
Публикация / подписка fanout-обменник, у каждого подписчика своя очередь
Маршрутизация direct по ключу
Темы topic по шаблону
RPC очередь ответа + reply_to + correlation_id
Отложенная доставка и повторы DLX + TTL (очередь ожидания возвращает в основную), плагин delayed-message
Повторы с ограничением счётчик в заголовках или x-delivery-limit у quorum-очередей, затем DLQ
Приоритеты x-max-priority (используйте осторожно)

Python (pika)

import json, pika

conn = pika.BlockingConnection(pika.ConnectionParameters('localhost', credentials=pika.PlainCredentials('app', 'secret')))
ch = conn.channel()
ch.confirm_delivery()
ch.exchange_declare('orders', exchange_type='topic', durable=True)
ch.queue_declare('orders.created', durable=True, arguments={
    'x-queue-type': 'quorum',
    'x-dead-letter-exchange': 'dlx',
})
ch.queue_bind('orders.created', 'orders', routing_key='orders.created')

ch.basic_publish('orders', 'orders.created', json.dumps({'id': 42}),
                 properties=pika.BasicProperties(delivery_mode=2, content_type='application/json', message_id='42'))

def handle(ch, method, props, body):
    try:
        process(json.loads(body))
        ch.basic_ack(method.delivery_tag)
    except Exception:
        ch.basic_nack(method.delivery_tag, requeue=False)      # в DLX

ch.basic_qos(prefetch_count=10)
ch.basic_consume('orders.created', handle)
ch.start_consuming()

Node.js (amqplib)

const amqp = require('amqplib');
const conn = await amqp.connect('amqp://app:secret@localhost');
const ch = await conn.createConfirmChannel();
await ch.assertQueue('tasks', { durable: true, arguments: { 'x-queue-type': 'quorum' } });
ch.sendToQueue('tasks', Buffer.from(JSON.stringify({ id: 1 })), { persistent: true });
await ch.prefetch(5);
await ch.consume('tasks', async (msg) => {
  try { await work(JSON.parse(msg.content)); ch.ack(msg); }
  catch { ch.nack(msg, false, false); }
});

Другие клиенты: Java (amqp-client, Spring AMQP @RabbitListener), .NET (RabbitMQ.Client, MassTransit), Go (amqp091-go), PHP (php-amqplib).

rabbitmqctl и плагины

rabbitmqctl status                                 rabbitmqctl cluster_status               rabbitmq-diagnostics check_running
rabbitmqctl list_queues name messages consumers state          rabbitmqctl list_queues -p / name messages_ready messages_unacknowledged
rabbitmqctl list_exchanges name type                rabbitmqctl list_bindings                rabbitmqctl list_connections name state
rabbitmqctl list_channels                           rabbitmqctl list_consumers
rabbitmqctl add_user app secret                     rabbitmqctl set_user_tags app monitoring
rabbitmqctl add_vhost shop                          rabbitmqctl set_permissions -p shop app ".*" ".*" ".*"      # configure, write, read
rabbitmqctl delete_user guest                        rabbitmqctl purge_queue <name> -p shop
rabbitmqctl list_policies                            rabbitmqctl set_policy ...                                    rabbitmqctl clear_policy name
rabbitmq-plugins list                                rabbitmq-plugins enable rabbitmq_management rabbitmq_prometheus rabbitmq_delayed_message_exchange
rabbitmqadmin declare queue name=q durable=true     # CLI поверх HTTP API

Порты: 5672 (AMQP), 5671 (AMQP+TLS), 15672 (веб-панель и HTTP API), 15692 (Prometheus), 25672 (межузловое), 4369 (epmd), 1883 (MQTT), 61613 (STOMP). Пользователь guest по умолчанию доступен только с localhost: заводите своих.

Кластер и отказоустойчивость

  • Узлы кластера: нечётное число (3, 5) для кворума. Quorum-очереди и стримы реплицируются по Raft.
  • Балансировщик (HAProxy, DNS, клиентские списки адресов) перед узлами; клиенты умеют переподключаться.
  • Разделение сети (partition): стратегии pause_minority и другие (cluster_partition_handling).
  • Мониторинг диска: при нехватке места срабатывает disk alarm, при памяти memory alarm (vm_memory_high_watermark): публикации блокируются.
  • Обновления: плавно, узел за узлом, с учётом совместимости версий.

Безопасность

TLS для клиентов и межузловой связи, отдельные учётные записи с минимальными правами на vhost, внешняя аутентификация (LDAP, OAuth 2.0), отключение guest, закрытие 15672 от интернета, ограничение размера сообщений.

Диагностика

Симптом Причина
Растёт messages_ready потребители не успевают или не подключены
Растёт messages_unacknowledged потребитель не отправляет ack, слишком большой prefetch
Сообщения «пропали» нет durable / persistent, нет привязки (обменник не знает, куда), mandatory не включён, очередь не создана
PRECONDITION_FAILED пересоздание очереди с другими аргументами: удалите или используйте политики
ACCESS_REFUSED права на vhost, неверный пароль
Блокировка публикаций memory или disk alarm
Дубли повторная доставка после сбоя: идемпотентность
«Ядовитые» сообщения крутятся бесконечно requeue=false + DLQ, x-delivery-limit

Включите mandatory и обработчик возврата (basic.return), чтобы видеть сообщения без маршрута.

Мониторинг

Веб-панель (:15672), плагин rabbitmq_prometheus + Grafana (готовые дашборды): глубина очередей, скорость публикации и доставки, число соединений, алерты на unacked, рост очередей, память и диск.

Kafka или RabbitMQ

RabbitMQ: сложная маршрутизация, очереди задач, приоритеты и TTL, подтверждения по каждому сообщению, низкая задержка. Kafka: поток событий, повторное чтение, огромная пропускная способность, хранение истории. Stream-очереди RabbitMQ закрывают часть сценариев Kafka. Подробнее см. шпаргалку Kafka.