RabbitMQ

Архитектура взаимодействия STOMP с RabbitMQ

RabbitMQ реализует поддержку STOMP через отдельный плагин, который преобразует STOMP-фреймы в внутреннюю модель брокера AMQP. В результате STOMP-клиенты получают доступ к очередям, обменникам и маршрутизации сообщений без прямой работы с AMQP-протоколом.

Ключевая особенность заключается в том, что STOMP выступает как текстовый протокол поверх TCP/WebSocket, а RabbitMQ выполняет роль брокера, обеспечивающего:

  • маршрутизацию сообщений через exchange;
  • хранение сообщений в очередях;
  • подтверждение доставки (ack/nack);
  • управление долговечностью сообщений;
  • распределение нагрузки между потребителями.

При использовании STOMP.js взаимодействие чаще всего происходит через WebSocket-соединение с RabbitMQ STOMP plugin.


Активация STOMP-плагина в RabbitMQ

Перед установкой соединения необходимо включить поддержку STOMP:

rabbitmq-plugins enable rabbitmq_stomp
rabbitmq-plugins enable rabbitmq_web_stomp

Первый плагин открывает TCP STOMP endpoint, второй — WebSocket-обёртку, необходимую для браузерных клиентов.

После активации становятся доступны порты:

  • 61613 — STOMP TCP
  • 15674 — STOMP over WebSocket

Подключение STOMP.js к RabbitMQ через WebSocket

STOMP.js использует WebSocket как транспорт, поверх которого формируется STOMP-сессия.

import { Client } from "@stomp/stompjs";

const client = new Client({
  brokerURL: "ws://localhost:15674/ws",
  reconnectDelay: 5000,
  heartbeatIncoming: 4000,
  heartbeatOutgoing: 4000,
});

При использовании RabbitMQ часто применяется логин и пароль:

client.connectHeaders = {
  login: "guest",
  passcode: "guest",
};

Важно учитывать, что учетная запись guest по умолчанию ограничена локальным доступом.


Формирование STOMP-сессии

После подключения создаётся сессия, в рамках которой доступны подписки и отправка сообщений.

client.onConn ect = () => {
  console.log("connected");
};

client.activate();

При установлении соединения RabbitMQ создаёт канал доставки сообщений, связанный с виртуальным хостом (vhost).


Модель маршрутизации RabbitMQ в STOMP

RabbitMQ не использует STOMP-дестинации напрямую. Вместо этого применяется трансляция:

Очереди

/queue/<queue_name>

Сообщения отправляются в конкретную очередь через default exchange.

client.publish({
  destination: "/queue/orders",
  body: JSON.stringify({ id: 1, status: "new" }),
});

Топики

/topic/<routing_key>

Используется fanout или topic exchange в RabbitMQ.

client.subscribe("/topic/news", (message) => {
  const payload = JSON.parse(message.body);
});

Прямой обмен (direct exchange)

В RabbitMQ STOMP может использоваться через привязку routing key:

/exchange/<exchange_name>/<routing_key>
client.publish({
  destination: "/exchange/logs/error",
  body: "critical error",
});

Подписка на сообщения

Подписка создаёт consumer, связанный с очередью или binding’ом exchange.

const subscription = client.subscribe("/queue/tasks", (message) => {
  const data = JSON.parse(message.body);
});

Каждое сообщение приходит как STOMP frame, содержащий:

  • body
  • headers
  • ack-id (при включённом режиме подтверждения)

Подтверждение доставки сообщений (ACK)

RabbitMQ поддерживает режимы подтверждения:

AUTO ACK

Сообщение считается доставленным сразу после получения.

client.subscribe("/queue/tasks", handler);

CLIENT ACK

Требуется явное подтверждение:

client.subscribe("/queue/tasks", (message) => {
  const data = JSON.parse(message.body);

  message.ack();
});

CLIENT INDIVIDUAL ACK

Подтверждение каждого сообщения отдельно:

message.ack({ mode: "client-individual" });

Отсутствие ACK приводит к повторной доставке сообщения.


Отправка сообщений в RabbitMQ через STOMP

Отправка формируется через publish:

client.publish({
  destination: "/queue/process",
  body: JSON.stringify({
    task: "resize-image",
    payload: {
      width: 300,
      height: 300,
    },
  }),
  headers: {
    persistent: "true",
  },
});

Основные заголовки RabbitMQ STOMP

  • persistent — сохранение сообщения на диск
  • priority — приоритет обработки
  • expiration — TTL сообщения
  • content-type — тип данных

TTL и истечение сообщений

RabbitMQ поддерживает время жизни сообщений через заголовок expiration.

client.publish({
  destination: "/queue/cache",
  body: "temp data",
  headers: {
    expiration: "10000",
  },
});

Значение указывается в миллисекундах, после чего сообщение удаляется из очереди.


Очереди и долговечность

Для устойчивой доставки важно согласование трёх параметров:

  • durable queue (создаётся на стороне RabbitMQ)
  • persistent message
  • durable exchange

STOMP.js не создаёт очередь напрямую — она должна существовать в RabbitMQ заранее или быть создана через административный API.


Обработка переподключений

RabbitMQ и WebSocket-соединения подвержены разрывам. STOMP.js поддерживает автоматическое восстановление:

const client = new Client({
  brokerURL: "ws://localhost:15674/ws",
  reconnectDelay: 3000,
});

После восстановления соединения подписки могут требовать повторной регистрации, в зависимости от конфигурации STOMP.js.


Heartbeat и контроль соединения

Heartbeat предотвращает разрыв idle-соединений.

const client = new Client({
  heartbeatIncoming: 10000,
  heartbeatOutgoing: 10000,
});

RabbitMQ проверяет активность канала и закрывает неактивные соединения.


Работа с виртуальными хостами (vhost)

RabbitMQ использует vhost для изоляции ресурсов.

Подключение через STOMP может включать vhost в URL:

ws://localhost:15674/ws

Или через заголовки:

client.connectHeaders = {
  login: "user",
  passcode: "pass",
  host: "/",
};

Ограничение скорости и prefetch

Хотя STOMP не управляет QoS напрямую, RabbitMQ применяет prefetch на уровне канала.

При большом потоке сообщений:

  • сообщения могут приходить быстрее, чем обрабатываются;
  • возможна перегрузка клиента;
  • необходимо учитывать backpressure на уровне обработки.

Очереди с несколькими подписчиками

RabbitMQ распределяет сообщения между подписчиками:

  • round-robin распределение;
  • конкурирующие consumers;
  • гарантированная доставка одному consumer при стандартной очереди.

STOMP.js не контролирует стратегию распределения — она задаётся RabbitMQ.


Ошибки соединения и диагностика

Типичные сценарии:

Connection refused

  • STOMP plugin не активирован
  • неверный порт WebSocket

404 destination

  • отсутствует очередь или binding
  • неверный exchange mapping

Access denied

  • ограничения пользователя RabbitMQ
  • неправильный vhost

Формат STOMP-фреймов в RabbitMQ

Сообщения передаются как текстовые STOMP frames:

SEND
destination:/queue/test
content-type:application/json

{"id":1}

Ответы:

MESSAGE
subscription:sub-0
message-id:123

{"id":1}

Безопасность соединений (WS vs WSS)

В production-средах используется TLS:

wss://broker.example.com:15674/ws

Особенности:

  • обязательная настройка сертификатов RabbitMQ;
  • корректная конфигурация reverse proxy (nginx, traefik);
  • защита от перехвата сообщений.

Масштабирование STOMP-клиентов

При росте нагрузки учитываются:

  • количество WebSocket-соединений;
  • лимиты RabbitMQ на file descriptors;
  • балансировка через reverse proxy;
  • использование нескольких vhost или кластеров RabbitMQ.

Интеграция с обменниками RabbitMQ

STOMP-дестинации транслируются в exchange bindings:

  • default exchange (direct)
  • fanout exchange (broadcast)
  • topic exchange (pattern routing)

Пример topic routing:

client.subscribe("/topic/orders.*", handler);

Со стороны RabbitMQ:

routing key: orders.created
routing key: orders.deleted

Управление жизненным циклом сообщений

RabbitMQ обеспечивает:

  • подтверждение доставки;
  • повторную постановку в очередь при отказе;
  • dead-letter exchange при ошибках;
  • TTL на уровне очереди и сообщения.

STOMP.js отражает только клиентскую часть этого процесса через ACK/NACK и headers.