🐰 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-сообщение.
Надёжность
- Подтверждения публикации (publisher confirms): брокер подтверждает приём.
- Подтверждения потребителя (
ack): сообщение удаляется послеbasic_ack;nack/rejectсrequeue=true|falseвозвращает или отбрасывает. Автоподтверждение (auto_ack) может терять сообщения. - Prefetch (
basic_qos): сколько неподтверждённых сообщений держит потребитель (1 для равномерного распределения тяжёлых задач, 10–100 для быстрых). - Обработка идемпотентна: возможны повторные доставки (
redelivered=true). - Dead Letter Exchange (DLX): куда попадают отклонённые, просроченные и переполнившие очередь сообщения.
- TTL: время жизни сообщения (
x-message-ttl) или очереди (x-expires). - Лимиты:
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.