Публикация в канал

Публикация сообщений — центральный механизм взаимодействия клиента с брокером сообщений в протоколе STOMP. В библиотеке STOMP.js отправка данных выполняется через метод publish(), который формирует STOMP-фрейм SEND и передаёт его брокеру через активное WebSocket-соединение.

Публикация используется для:

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

Метод publish()

Базовая публикация сообщения выглядит следующим образом:

client.publish({
    destination: '/topic/chat',
    body: 'Привет'
});

После вызова метода библиотека:

  1. Формирует STOMP-фрейм SEND.
  2. Добавляет заголовки.
  3. Кодирует тело сообщения.
  4. Передаёт данные через WebSocket.
  5. Отправляет пакет брокеру.

Структура публикации

Метод publish() принимает объект конфигурации.

Основные параметры:

Параметр Назначение
destination Адрес канала
body Тело сообщения
headers STOMP-заголовки
binaryBody Бинарные данные
skipContentLengthHeader Отключение заголовка content-length

destination

Поле destination определяет адрес назначения сообщения.

Примеры:

destination: '/topic/chat'
destination: '/queue/tasks'
destination: '/app/sendMessage'

Тип адреса зависит от брокера:

Префикс Назначение
/topic/ Публикация для группы подписчиков
/queue/ Очередь сообщений
/app/ Серверный endpoint
/exchange/ RabbitMQ exchange
/user/ Персональные сообщения

Простая публикация текста

Наиболее распространённый вариант:

client.publish({
    destination: '/topic/news',
    body: 'Новое сообщение'
});

Тело автоматически отправляется как строка.


Публикация JSON

Чаще всего STOMP используется для обмена JSON-структурами.

const message = {
    id: 15,
    text: 'Привет',
    author: 'Alex'
};

client.publish({
    destination: '/topic/chat',
    body: JSON.stringify(message)
});

На сервере данные обычно десериализуются обратно в объект.


content-type

При отправке JSON рекомендуется указывать MIME-тип.

client.publish({
    destination: '/topic/chat',
    body: JSON.stringify({
        text: 'Сообщение'
    }),
    headers: {
        'content-type': 'application/json'
    }
});

Популярные типы:

MIME Назначение
text/plain Обычный текст
application/json JSON
application/xml XML
application/octet-stream Бинарные данные

STOMP-заголовки

Публикация может содержать произвольные заголовки.

client.publish({
    destination: '/topic/events',
    body: 'event',
    headers: {
        priority: 'high',
        type: 'notification'
    }
});

Заголовки используются для:

  • маршрутизации;
  • фильтрации;
  • авторизации;
  • идентификации сообщений;
  • TTL;
  • приоритетов;
  • корреляции запросов.

Пользовательские идентификаторы сообщений

Нередко клиент самостоятельно назначает идентификатор:

client.publish({
    destination: '/topic/chat',
    body: 'Сообщение',
    headers: {
        messageId: crypto.randomUUID()
    }
});

Это полезно для:

  • дедупликации;
  • отслеживания доставки;
  • логирования;
  • корреляции ответов.

Отправка временных меток

client.publish({
    destination: '/topic/logs',
    body: JSON.stringify({
        message: 'System started',
        timestamp: Date.now()
    })
});

В распределённых системах временные метки особенно важны для:

  • аудита;
  • анализа задержек;
  • восстановления последовательности событий.

Публикация после подключения

Отправка сообщений возможна только после успешного соединения.

Корректный вариант:

client.onConn ect = () => {
    client.publish({
        destination: '/topic/chat',
        body: 'Подключение успешно'
    });
};

client.activate();

Некорректный вариант:

client.activate();

client.publish({
    destination: '/topic/chat',
    body: 'Ошибка'
});

Во втором случае соединение может ещё не успеть установиться.


Проверка состояния соединения

if (client.connected) {
    client.publish({
        destination: '/topic/chat',
        body: 'Сообщение'
    });
}

Свойство connected показывает текущее состояние клиента.


Публикация в обработчике событий

button.addEventListener('click', () => {
    client.publish({
        destination: '/app/message',
        body: JSON.stringify({
            text: input.value
        })
    });
});

Такой подход используется в:

  • чатах;
  • совместных редакторах;
  • игровых приложениях;
  • системах уведомлений.

Массовая отправка сообщений

for (let i = 0; i < 100; i++) {
    client.publish({
        destination: '/queue/tasks',
        body: JSON.stringify({
            taskId: i
        })
    });
}

При высокой нагрузке необходимо учитывать:

  • пропускную способность WebSocket;
  • ограничения брокера;
  • размер очередей;
  • объём памяти;
  • скорость обработки подписчиков.

Публикация с подтверждением receipt

STOMP поддерживает подтверждение обработки команды брокером.

client.publish({
    destination: '/topic/chat',
    body: 'Сообщение',
    headers: {
        receipt: 'msg-001'
    }
});

Обработка receipt:

client.onRece ipt = (frame) => {
    console.log('Получено подтверждение');
    console.log(frame.headers['receipt-id']);
};

Receipt подтверждает:

  • получение команды брокером;
  • регистрацию сообщения;
  • успешную обработку STOMP-фрейма.

Receipt не гарантирует доставку подписчику.


Использование receipt для критически важных сообщений

client.publish({
    destination: '/queue/payments',
    body: JSON.stringify({
        amount: 100
    }),
    headers: {
        receipt: 'payment-100'
    }
});

Подобный механизм используется в:

  • финансовых системах;
  • обработке заказов;
  • очередях задач;
  • критически важных уведомлениях.

Работа с бинарными данными

STOMP.js поддерживает передачу бинарных данных через binaryBody.

const bytes = new Uint8Array([1, 2, 3, 4]);

client.publish({
    destination: '/topic/binary',
    binaryBody: bytes,
    headers: {
        'content-type': 'application/octet-stream'
    }
});

Поддерживаются:

  • Uint8Array;
  • ArrayBuffer;
  • бинарные буферы.

Передача файлов

Пример отправки файла:

fileInput.addEventListener('change', async (event) => {
    const file = event.target.files[0];

    const arrayBuffer = await file.arrayBuffer();

    client.publish({
        destination: '/topic/files',
        binaryBody: new Uint8Array(arrayBuffer),
        headers: {
            filename: file.name,
            'content-type': file.type
        }
    });
});

Размер сообщений

Брокеры обычно ограничивают максимальный размер пакета.

Типичные ограничения:

Брокер Ограничение
RabbitMQ configurable
ActiveMQ configurable
Apollo configurable

Слишком большие сообщения могут:

  • разрывать соединение;
  • вызывать ошибки памяти;
  • блокировать обработку;
  • увеличивать задержки.

Фрагментация больших данных

Вместо одного огромного сообщения часто используют разбиение на части.

client.publish({
    destination: '/topic/upload',
    body: JSON.stringify({
        fileId: 'abc',
        chunk: 1,
        total: 10,
        data: chunkData
    })
});

Публикация без content-length

По умолчанию STOMP.js добавляет заголовок content-length.

Иногда его требуется отключить.

client.publish({
    destination: '/topic/chat',
    body: 'Hello',
    skipContentLengthHeader: true
});

Это может понадобиться при несовместимости со старым брокером.


Публикация XML

client.publish({
    destination: '/topic/xml',
    body: `
        <message>
            <text>Hello</text>
        </message>
    `,
    headers: {
        'content-type': 'application/xml'
    }
});

Публикация с авторизационными данными

Иногда токен передаётся в заголовках сообщения.

client.publish({
    destination: '/app/private',
    body: 'secret',
    headers: {
        Authorization: 'Bearer token'
    }
});

Однако чаще авторизация выполняется на уровне CONNECT-фрейма.


Отправка сообщений в Spring

При интеграции с Spring Framework часто используются endpoint’ы:

client.publish({
    destination: '/app/chat.send',
    body: JSON.stringify({
        text: 'Hello'
    })
});

Серверный обработчик:

@MessageMapping("/chat.send")
public void send(ChatMessage message) {
    // обработка
}

Работа с RabbitMQ

При использовании RabbitMQ возможны специальные адреса:

client.publish({
    destination: '/exchange/chat/messages',
    body: 'Hello RabbitMQ'
});

Также поддерживаются:

/queue/
/topic/
/exchange/
/amq/queue/

Обработка ошибок публикации

Сам вызов publish() обычно не выбрасывает исключение доставки.

Но возможны ошибки:

  • отсутствующее соединение;
  • закрытый WebSocket;
  • ошибка сериализации;
  • превышение лимита;
  • отказ брокера.

Пример проверки:

try {
    if (!client.connected) {
        throw new Error('Нет соединения');
    }

    client.publish({
        destination: '/topic/chat',
        body: 'Hello'
    });
} catch (error) {
    console.error(error);
}

Буферизация сообщений

Некоторые приложения временно сохраняют сообщения до подключения.

const pendingMessages = [];

function send(message) {
    if (client.connected) {
        client.publish({
            destination: '/topic/chat',
            body: message
        });
    } else {
        pendingMessages.push(message);
    }
}

После подключения:

client.onConn ect = () => {
    pendingMessages.forEach(message => {
        client.publish({
            destination: '/topic/chat',
            body: message
        });
    });

    pendingMessages.length = 0;
};

Повторная отправка

При сетевых сбоях используется retry-механизм.

function publishWithRetry(message, retries = 3) {
    try {
        client.publish({
            destination: '/topic/chat',
            body: message
        });
    } catch (error) {
        if (retries > 0) {
            setTimeout(() => {
                publishWithRetry(message, retries - 1);
            }, 1000);
        }
    }
}

Производительность публикации

На производительность влияют:

  • размер сообщений;
  • количество заголовков;
  • скорость сериализации JSON;
  • нагрузка брокера;
  • частота отправки;
  • качество сети.

Минимизация нагрузки

Для оптимизации применяются:

  • batching;
  • сжатие данных;
  • бинарные форматы;
  • ограничение частоты отправки;
  • агрегация событий.

Batch-публикация

const events = [];

for (let i = 0; i < 100; i++) {
    events.push({
        id: i
    });
}

client.publish({
    destination: '/topic/batch',
    body: JSON.stringify(events)
});

Такой подход снижает:

  • число STOMP-фреймов;
  • нагрузку на сеть;
  • количество операций сериализации.

Публикация heartbeat-событий

Некоторые приложения публикуют собственные heartbeat-сообщения.

setInterval(() => {
    client.publish({
        destination: '/topic/heartbeat',
        body: JSON.stringify({
            timestamp: Date.now()
        })
    });
}, 5000);

Публикация событий приложения

STOMP часто используется как event bus.

client.publish({
    destination: '/topic/events',
    body: JSON.stringify({
        type: 'USER_CREATED',
        payload: {
            id: 10
        }
    })
});

Стандартная структура событий обычно включает:

Поле Назначение
type Тип события
payload Данные
timestamp Время
version Версия схемы
source Источник события

Идемпотентность сообщений

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

Для этого используются:

client.publish({
    destination: '/queue/orders',
    body: JSON.stringify({
        operationId: crypto.randomUUID(),
        amount: 500
    })
});

Идемпотентность особенно важна для:

  • платежей;
  • создания заказов;
  • изменения баланса;
  • резервирования ресурсов.

Логирование публикаций

function publish(destination, body) {
    console.log('SEND', destination, body);

    client.publish({
        destination,
        body: JSON.stringify(body)
    });
}

Логирование помогает:

  • анализировать трафик;
  • искать ошибки;
  • восстанавливать последовательность событий;
  • отслеживать нагрузку.

Обёртка над publish()

В крупных проектах публикацию часто инкапсулируют.

class MessageBus {
    constructor(client) {
        this.client = client;
    }

    send(destination, payload) {
        this.client.publish({
            destination,
            body: JSON.stringify(payload),
            headers: {
                'content-type': 'application/json'
            }
        });
    }
}

Преимущества:

  • единая сериализация;
  • централизованные заголовки;
  • логирование;
  • retry;
  • мониторинг;
  • унификация API.

Типизация сообщений в TypeScript

interface ChatMessage {
    id: number;
    text: string;
}

function publishChatMessage(message: ChatMessage) {
    client.publish({
        destination: '/topic/chat',
        body: JSON.stringify(message)
    });
}

Типизация уменьшает количество ошибок при сериализации и обработке событий.


Частые ошибки при публикации

Отправка до подключения

client.publish({
    destination: '/topic/chat',
    body: 'Ошибка'
});

Отсутствие сериализации JSON

Неправильно:

body: {
    text: 'Hello'
}

Правильно:

body: JSON.stringify({
    text: 'Hello'
})

Неверный destination

destination: 'chat'

Корректный вариант:

destination: '/topic/chat'

Отсутствие content-type

headers: {
    'content-type': 'application/json'
}

Архитектурные особенности публикации

В больших системах публикация сообщений обычно строится вокруг:

  • event-driven architecture;
  • message bus;
  • CQRS;
  • микросервисов;
  • асинхронного взаимодействия;
  • очередей задач.

STOMP.js при этом выступает клиентским уровнем транспорта между браузером и брокером сообщений.