HornetQ

HornetQ представляет собой высокопроизводительный брокер сообщений, построенный с ориентацией на масштабируемые распределённые системы и асинхронный обмен данными. Его архитектура изначально проектировалась как часть экосистемы JBoss и включает поддержку нескольких протоколов поверх единого ядра маршрутизации сообщений. Одним из ключевых протоколов взаимодействия выступает STOMP, что делает интеграцию с JavaScript-клиентами через STOMP.js прямолинейной и предсказуемой.

HornetQ реализует модель брокера сообщений с разделением на адреса (addresses) и очереди (queues). Адрес выступает логической точкой публикации, а очередь — конкретным потребителем сообщений. При использовании STOMP поверх HornetQ происходит трансляция STOMP-команд в внутренние операции брокера.

Основные компоненты:

  • STOMP acceptor — сетевой endpoint, принимающий STOMP-соединения
  • Address registry — реестр маршрутов сообщений
  • Queue binding layer — связывание адресов с очередями
  • Persistence layer — хранение сообщений (опционально)
  • Delivery scheduler — управление доставкой и retry

STOMP-клиент, подключающийся через WebSocket или TCP, не взаимодействует напрямую с очередями HornetQ. Он работает с абстракцией destination, которая на стороне брокера преобразуется в address/queue mapping.

Включение STOMP в HornetQ

Конфигурация STOMP acceptor обычно задаётся в hornetq-configuration.xml. Пример логической конфигурации:

<acceptors>
   <acceptor name="stomp-acceptor">
      tcp://0.0.0.0:61613?protocols=STOMP
   </acceptor>
</acceptors>

В некоторых конфигурациях используется мультипротокольный acceptor:

<acceptor name="netty">
   tcp://0.0.0.0:61616?protocols=CORE,STOMP,AMQP
</acceptor>

Ключевым моментом является параметр protocols=STOMP, который активирует STOMP codec внутри Netty pipeline.

STOMP модель взаимодействия

STOMP (Simple Text Oriented Messaging Protocol) представляет собой текстовый протокол поверх TCP/WebSocket. Основные операции:

  • CONNECT — установление соединения
  • SEND — отправка сообщения
  • SUBSCRIBE — подписка на очередь/топик
  • UNSUBSCRIBE — отмена подписки
  • ACK / NACK — подтверждение доставки
  • DISCONNECT — завершение сессии

HornetQ интерпретирует destination следующим образом:

  • /queue/name → point-to-point очередь
  • /topic/name → pub-sub модель

Подключение STOMP.js к HornetQ

STOMP.js используется как клиентская реализация STOMP поверх WebSocket или TCP-over-proxy.

Базовое подключение:

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

const client = new Client({
  brokerURL: 'ws://localhost:61614/stomp',
  reconnectDelay: 5000,
  debug: (msg) => console.log(msg)
});

client.onConn ect = (frame) => {
  console.log('Connected:', frame.headers);

  client.subscribe('/queue/test', (message) => {
    const body = message.body;
    console.log('Received:', body);
  });
};

client.activate();

В случае HornetQ важно, чтобы был включён WebSocket STOMP endpoint (часто через proxy или servlet container).

Отправка сообщений в очередь

client.publish({
  destination: '/queue/test',
  body: JSON.stringify({
    type: 'event',
    payload: {
      id: 123,
      status: 'ok'
    }
  })
});

HornetQ маршрутизирует сообщение в очередь test, создавая binding между address /queue/test и соответствующей queue instance.

Подписка и модели доставки

STOMP.js поддерживает несколько режимов подписки, которые напрямую влияют на поведение HornetQ:

AUTO ACK

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

client.subscribe('/queue/test', (message) => {
  console.log(message.body);
}, { ack: 'auto' });

CLIENT ACK

Ручное подтверждение доставки:

client.subscribe('/queue/test', (message) => {
  console.log(message.body);
  message.ack();
}, { ack: 'client' });

HornetQ удерживает сообщение в unacknowledged state до получения ACK-фрейма.

CLIENT INDIVIDUAL ACK

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

client.subscribe('/queue/test', (message) => {
  message.ack();
}, { ack: 'client-individual' });

Транзакции и гарантии доставки

HornetQ поддерживает транзакционные сессии STOMP. В STOMP.js это выражается через headers:

client.begin('tx1');

client.publish({
  destination: '/queue/test',
  body: 'message in transaction',
  headers: { transaction: 'tx1' }
});

client.commit('tx1');

При rollback сообщение не фиксируется в очереди.

Heartbeat и стабильность соединения

STOMP поддерживает heartbeat механизм для контроля соединения:

const client = new Client({
  brokerURL: 'ws://localhost:61614/stomp',
  heartbeatIncoming: 10000,
  heartbeatOutgoing: 10000
});

HornetQ использует heartbeat для очистки зависших сессий и освобождения consumer lock.

Маршрутизация сообщений внутри HornetQ

После получения STOMP SEND:

  1. Parse destination
  2. Resolve address mapping
  3. Apply routing type (MULTICAST / ANYCAST)
  4. Select queue(s)
  5. Persist message (если включено)
  6. Dispatch to consumers

MULTICAST соответствует topic-модели, ANYCAST — queue-модели.

Приоритеты сообщений

HornetQ поддерживает приоритеты через STOMP headers:

client.publish({
  destination: '/queue/test',
  body: 'high priority message',
  headers: {
    priority: 9
  }
});

Приоритет влияет на порядок доставки внутри очереди, но не гарантирует абсолютную сортировку при высокой конкуренции consumers.

TTL сообщений

TTL задаётся через header:

client.publish({
  destination: '/queue/test',
  body: 'temporary message',
  headers: {
    'expires': Date.now() + 60000
  }
});

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

Масштабирование потребителей

STOMP-подписчики на одну очередь в HornetQ работают в режиме конкурирующих consumers. Сообщение доставляется только одному из подписчиков:

  • load balancing между consumers
  • распределение через round-robin или weighted dispatch
  • back-pressure при медленных consumers

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

Типичные STOMP ошибки:

  • message: malformed frame — неправильный формат STOMP frame
  • access denied — отсутствие прав на destination
  • no route — отсутствие binding между address и queue
  • consumer closed — разрыв подписки

STOMP.js позволяет отслеживать ошибки:

client.onStompEr ror = (frame) => {
  console.error('Broker error:', frame.body);
};

Буферизация и производительность

STOMP.js использует внутренний буфер outbound frames. При высокой нагрузке:

  • сообщения могут накапливаться в memory buffer
  • требуется tuning maxWebSocketFrameSize на стороне брокера
  • важно контролировать burst SEND операций

HornetQ, в свою очередь, оптимизирует throughput через batching и asynchronous IO.

Сценарии интеграции

Реалтайм уведомления

  • WebSocket STOMP client
  • HornetQ topic /topic/notifications
  • push-архитектура без polling

Очередь задач

  • /queue/jobs
  • worker consumers на backend
  • гарантированная доставка

Event-driven архитектура

  • domain events через topics
  • multiple subscribers
  • eventual consistency между сервисами

Внутреннее преобразование STOMP в HornetQ Core

STOMP frame:

SEND
destination:/queue/test
content-type:text/plain

hello

преобразуется в internal message:

  • Address: queue/test
  • RoutingType: ANYCAST
  • Body: “hello”
  • Properties: headers map
  • DeliveryCount: 0

Далее сообщение попадает в paging store при необходимости.

Безопасность соединений

HornetQ поддерживает:

  • simple login/password authentication
  • JAAS integration
  • SSL/TLS поверх STOMP TCP или WebSocket

STOMP.js передаёт credentials:

client.connectHeaders = {
  login: 'user',
  passcode: 'password'
};

Особенности взаимодействия с WebSocket транспортом

При использовании WebSocket:

  • STOMP frame encapsulated in WebSocket messages
  • no persistent TCP session exposure
  • easier NAT traversal
  • requires broker-side WebSocket acceptor or proxy layer

HornetQ в классической конфигурации чаще использует TCP STOMP, WebSocket добавляется через внешние компоненты.

Ограничения модели

  • отсутствие сложной маршрутизации как в ESB
  • ограниченная трансформация сообщений
  • зависимость от конфигурации acceptor
  • необходимость ручного управления ack/transactions при высокой надёжности

Поведение при отказах

При разрыве соединения:

  • неподтверждённые сообщения возвращаются в очередь
  • consumers пересоздаются при reconnect
  • STOMP.js reconnectDelay автоматически восстанавливает сессию
  • HornetQ сохраняет состояние очередей в persistence store при включённой персистентности