Обработка входящих сообщений

После установки соединения и оформления подписки STOMP-клиент начинает получать сообщения от брокера. В STOMP.js обработка входящих данных строится вокруг callback-функций, передаваемых в метод subscribe(), а также обработчиков системных событий клиента.

Каждое входящее сообщение поступает в виде объекта IMessage, содержащего:

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

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

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

const client = new Client({
    brokerURL: 'ws://localhost:15674/ws'
});

client.onConn ect = () => {

    client.subscribe('/topic/news', (message) => {

        console.log('Получено сообщение');
        console.log(message.body);

    });

};

client.activate();

После получения данных callback вызывается автоматически.


Структура объекта IMessage

В callback подписки передаётся объект сообщения:

client.subscribe('/topic/chat', (message) => {

    console.log(message);

});

Наиболее важные свойства:

Свойство Назначение
body Тело сообщения
headers Заголовки STOMP
ack() Подтверждение обработки
nack() Отказ от обработки
binaryBody Бинарные данные
command Тип STOMP-фрейма

Пример анализа содержимого:

client.subscribe('/topic/orders', (message) => {

    console.log(message.command);
    console.log(message.headers);
    console.log(message.body);

});

Получение текстового содержимого

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

client.subscribe('/topic/logs', (message) => {

    const text = message.body;

    console.log(text);

});

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

Если сервер отправил JSON, потребуется ручной разбор.


Обработка JSON-сообщений

Большинство приложений используют JSON как основной формат обмена.

Пример сообщения от сервера:

{
    "id": 15,
    "status": "done",
    "price": 4500
}

Разбор сообщения:

client.subscribe('/topic/orders', (message) => {

    const data = JSON.parse(message.body);

    console.log(data.id);
    console.log(data.status);
    console.log(data.price);

});

Проверка корректности JSON

Некорректный JSON вызывает исключение:

JSON.parse(message.body);

Поэтому обработку желательно защищать:

client.subscribe('/topic/orders', (message) => {

    try {

        const data = JSON.parse(message.body);

        console.log(data);

    } catch (error) {

        console.error('Ошибка JSON');
        console.error(error);

    }

});

Обработка разных типов сообщений

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

Например:

{
    "type": "created",
    "payload": {
        "id": 5
    }
}
{
    "type": "deleted",
    "payload": {
        "id": 5
    }
}

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

client.subscribe('/topic/events', (message) => {

    const data = JSON.parse(message.body);

    switch (data.type) {

        case 'created':

            console.log('Создание');
            console.log(data.payload);

            break;

        case 'deleted':

            console.log('Удаление');
            console.log(data.payload);

            break;

        default:

            console.log('Неизвестный тип');

    }

});

Использование заголовков сообщения

STOMP поддерживает передачу метаданных через headers.

Пример чтения заголовков:

client.subscribe('/topic/files', (message) => {

    console.log(message.headers);

});

Получение конкретного заголовка:

client.subscribe('/topic/files', (message) => {

    const contentType = message.headers['content-type'];

    console.log(contentType);

});

Content-Type входящих сообщений

Сервер часто указывает MIME-тип:

client.subscribe('/topic/data', (message) => {

    const type = message.headers['content-type'];

    if (type === 'application/json') {

        const data = JSON.parse(message.body);

        console.log(data);

    }

});

Обработка бинарных сообщений

STOMP.js поддерживает бинарные payload.

Пример:

client.subscribe('/topic/binary', (message) => {

    const bytes = message.binaryBody;

    console.log(bytes);

});

binaryBody содержит Uint8Array.


Преобразование бинарных данных в строку

client.subscribe('/topic/binary', (message) => {

    const decoder = new TextDecoder();

    const text = decoder.decode(message.binaryBody);

    console.log(text);

});

Преобразование бинарных данных в JSON

client.subscribe('/topic/binary', (message) => {

    const decoder = new TextDecoder();

    const text = decoder.decode(message.binaryBody);

    const data = JSON.parse(text);

    console.log(data);

});

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

По умолчанию большинство брокеров используют автоматическое подтверждение доставки.

В режиме client подтверждение выполняется вручную.

Подписка:

client.subscribe(
    '/queue/tasks',
    (message) => {

        console.log(message.body);

        message.ack();

    },
    {
        ack: 'client'
    }
);

Подтверждение после успешной обработки

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

client.subscribe(
    '/queue/tasks',
    async (message) => {

        try {

            const data = JSON.parse(message.body);

            await saveTask(data);

            message.ack();

        } catch (error) {

            console.error(error);

        }

    },
    {
        ack: 'client'
    }
);

Отказ от обработки через nack()

Если сообщение не удалось обработать:

client.subscribe(
    '/queue/tasks',
    (message) => {

        try {

            processTask(message.body);

            message.ack();

        } catch (error) {

            message.nack();

        }

    },
    {
        ack: 'client'
    }
);

Поведение после nack() зависит от брокера:

  • повторная доставка;
  • помещение в dead-letter queue;
  • удаление сообщения;
  • задержанный retry.

Использование async/await

STOMP.js корректно работает с асинхронными callback.

client.subscribe('/topic/users', async (message) => {

    const user = JSON.parse(message.body);

    await loadProfile(user.id);

    console.log('Профиль загружен');

});

Обработка ошибок внутри callback

Ошибки внутри callback не перехватываются автоматически клиентом.

Опасный вариант:

client.subscribe('/topic/data', (message) => {

    const data = JSON.parse(message.body);

    runOperation(data);

});

Без обработки исключение может нарушить логику приложения.

Безопасный вариант:

client.subscribe('/topic/data', (message) => {

    try {

        const data = JSON.parse(message.body);

        runOperation(data);

    } catch (error) {

        console.error(error);

    }

});

Централизованная обработка сообщений

В крупных приложениях callback обычно делегирует работу отдельным обработчикам.

Пример:

function handleOrder(message) {

    const data = JSON.parse(message.body);

    console.log(data);

}

client.subscribe('/topic/orders', handleOrder);

Разделение обработчиков по каналам

client.subscribe('/topic/users', handleUsers);

client.subscribe('/topic/orders', handleOrders);

client.subscribe('/topic/notifications', handleNotifications);

Такой подход:

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

Фильтрация сообщений

Иногда требуется игнорировать часть входящих данных.

client.subscribe('/topic/orders', (message) => {

    const data = JSON.parse(message.body);

    if (data.status !== 'active') {
        return;
    }

    console.log(data);

});

Дедупликация сообщений

Некоторые брокеры допускают повторную доставку.

Пример простой дедупликации:

const processed = new Set();

client.subscribe('/queue/tasks', (message) => {

    const id = message.headers['message-id'];

    if (processed.has(id)) {
        return;
    }

    processed.add(id);

    console.log(message.body);

});

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

При высокой нагрузке поток сообщений может перегружать интерфейс.

Пример debounce:

let timeout;

client.subscribe('/topic/search', (message) => {

    clearTimeout(timeout);

    timeout = setTimeout(() => {

        console.log(message.body);

    }, 300);

});

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

Иногда данные накапливаются перед пакетной обработкой.

const buffer = [];

client.subscribe('/topic/events', (message) => {

    buffer.push(JSON.parse(message.body));

});

Периодическая обработка:

setInterval(() => {

    if (buffer.length === 0) {
        return;
    }

    console.log(buffer);

    buffer.length = 0;

}, 5000);

Работа с большим количеством сообщений

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

  • скорость JSON.parse;
  • нагрузку на DOM;
  • количество логирования;
  • объём памяти;
  • длительность callback.

Плохой вариант:

client.subscribe('/topic/stream', (message) => {

    heavyOperation(message.body);

});

Если обработка медленная, очередь сообщений начинает расти.


Асинхронная разгрузка обработки

Иногда полезно выносить тяжёлые операции:

client.subscribe('/topic/stream', (message) => {

    queueMicrotask(() => {

        heavyOperation(message.body);

    });

});

Использование Web Worker

Для тяжёлых вычислений:

const worker = new Worker('worker.js');

client.subscribe('/topic/stats', (message) => {

    worker.postMessage(message.body);

});

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

Пример внутри эффекта:

useEffect(() => {

    const subscription = client.subscribe(
        '/topic/chat',
        (message) => {

            setMessages((prev) => [
                ...prev,
                JSON.parse(message.body)
            ]);

        }
    );

    return () => {
        subscription.unsubscribe();
    };

}, []);

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

onMounted(() => {

    subscription = client.subscribe(
        '/topic/chat',
        (message) => {

            messages.value.push(
                JSON.parse(message.body)
            );

        }
    );

});

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

this.client.subscribe(
    '/topic/chat',
    (message) => {

        const data = JSON.parse(message.body);

        this.messages.push(data);

    }
);

Защита от утечек памяти

Подписки необходимо удалять:

const subscription = client.subscribe(
    '/topic/chat',
    handler
);

subscription.unsubscribe();

Без отписки callback продолжит получать сообщения.


Логирование входящих сообщений

Для диагностики полезно включать логирование:

client.subscribe('/topic/debug', (message) => {

    console.log('HEADERS');
    console.log(message.headers);

    console.log('BODY');
    console.log(message.body);

});

Анализ необработанных STOMP-фреймов

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

const client = new Client({

    brokerURL: 'ws://localhost:15674/ws',

    debug(str) {

        console.log(str);

    }

});

Пример вывода:

<<< MESSAGE
subscription:sub-0
message-id:007
destination:/topic/chat
content-type:application/json

Работа с heartbeat во время получения сообщений

Heartbeat помогает обнаруживать обрыв соединения.

const client = new Client({

    brokerURL: 'ws://localhost:15674/ws',

    heartbeatIncoming: 4000,
    heartbeatOutgoing: 4000

});

Если heartbeat перестаёт поступать, клиент инициирует переподключение.


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

После реконнекта подписки создаются заново внутри onConnect.

client.onConn ect = () => {

    client.subscribe('/topic/chat', handler);

};

Если подписка была создана вне onConnect, после reconnect сообщения перестанут поступать.


Обработка системных ошибок STOMP

STOMP-сервер может отправлять ERROR-фреймы.

client.onStompEr ror = (frame) => {

    console.error(frame.headers['message']);

    console.error(frame.body);

};

Обработка ошибок WebSocket

client.onWebSocketEr ror = (error) => {

    console.error(error);

};

Обработка закрытия соединения

client.onWebSocketCl ose = (event) => {

    console.log(event.code);
    console.log(event.reason);

};

Типичная схема обработки входящего сообщения

Полный пример:

client.subscribe(
    '/queue/orders',
    async (message) => {

        try {

            const contentType =
                message.headers['content-type'];

            if (contentType !== 'application/json') {

                message.nack();

                return;

            }

            const data = JSON.parse(message.body);

            await processOrder(data);

            message.ack();

        } catch (error) {

            console.error(error);

            message.nack();

        }

    },
    {
        ack: 'client'
    }
);

Архитектура промышленной обработки сообщений

В production-системах обработка обычно включает:

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

STOMP.js предоставляет только транспортный уровень. Вся прикладная логика обработки реализуется на стороне приложения.