Методы подписки

Подписка в STOMP.js — механизм получения сообщений от брокера через подписанные назначения (destination). После успешного подключения клиент может подписываться на очереди, топики и пользовательские каналы. Подписка создаёт постоянный канал доставки сообщений от брокера к клиенту.

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

  1. Клиент подключается к брокеру.
  2. Клиент подписывается на нужный канал.
  3. Брокер начинает отправлять сообщения.
  4. Callback-функция получает данные.
  5. При необходимости подписка удаляется.

Основным методом подписки является subscribe().


Метод subscribe()

Базовый синтаксис

const subscription = client.subscribe(destination, callback);

Параметры

Параметр Описание
destination Адрес подписки
callback Функция обработки сообщений

Простейшая подписка

client.subscribe('/topic/chat', (message) => {
    console.log(message.body);
});

После получения сообщения callback получает объект IMessage.


Объект сообщения IMessage

Метод подписки всегда передаёт объект сообщения.

Основные свойства

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

    console.log(message.body);

    console.log(message.headers);

    console.log(message.command);

});

Свойство body

body содержит текстовое содержимое сообщения.

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

    console.log(message.body);

});

Работа с JSON

Практически всегда сообщения сериализуются в JSON.

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

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

    console.log(data.id);

    console.log(data.status);

});

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

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

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

    try {

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

        console.log(data);

    } catch (error) {

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

    }

});

Подписка с заголовками

Метод subscribe() поддерживает третий параметр — объект заголовков.

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

Режимы подтверждения сообщений

Параметр ack определяет способ подтверждения доставки.

Возможные значения

Значение Описание
auto Автоматическое подтверждение
client Ручное подтверждение
client-individual Индивидуальное подтверждение

Автоматическое подтверждение

Режим auto

По умолчанию используется auto.

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

        console.log(message.body);

    },
    {
        ack: 'auto'
    }
);

После доставки сообщение автоматически считается обработанным.


Ручное подтверждение

Режим client

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

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

        console.log(message.body);

        message.ack();

    },
    {
        ack: 'client'
    }
);

Если подтверждение не отправлено, брокер может повторно доставить сообщение.


Зачем нужен manual ack

Ручное подтверждение используется:

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

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

Если клиент отключился до подтверждения:

message.ack();

сообщение может вернуться обратно в очередь.

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


Режим client-individual

Индивидуальное подтверждение

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

        processTask(message);

        message.ack();

    },
    {
        ack: 'client-individual'
    }
);

В этом режиме каждое сообщение подтверждается независимо.


Отличие client от client-individual

client

Подтверждает все предыдущие сообщения одновременно.

client-individual

Подтверждает только текущее сообщение.


Метод nack()

Отрицательное подтверждение

Сообщение можно отклонить.

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

        try {

            processTask(message.body);

            message.ack();

        } catch (error) {

            message.nack();

        }

    },
    {
        ack: 'client-individual'
    }
);

Когда используется nack()

nack() применяется:

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

Объект подписки

Возвращаемое значение subscribe()

Метод возвращает объект подписки.

const subscription = client.subscribe(
    '/topic/news',
    (message) => {
        console.log(message.body);
    }
);

Структура подписки

Основные элементы:

Свойство Назначение
id Идентификатор подписки
unsubscribe() Отписка

Метод unsubscribe()

Удаление подписки

const subscription = client.subscribe(
    '/topic/news',
    (message) => {
        console.log(message.body);
    }
);

subscription.unsubscribe();

После вызова сообщения больше не поступают.


Важность unsubscribe()

Отсутствие отписки вызывает:

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

Повторные подписки

Ошибка множественных subscribe()

function init() {

    client.subscribe('/topic/chat', (message) => {
        console.log(message.body);
    });

}

Если функция вызывается несколько раз, количество подписок начинает расти.


Последствия дублирования

Проблемы:

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

Защита от повторных подписок

let subscription = null;

function initSubscription() {

    if (subscription) {
        return;
    }

    subscription = client.subscribe(
        '/topic/chat',
        (message) => {
            console.log(message.body);
        }
    );

}

Хранение подписок

Централизованное управление

При большом количестве каналов подписки обычно хранятся в коллекции.

const subscriptions = {};

Регистрация подписок

subscriptions.chat = client.subscribe(
    '/topic/chat',
    (message) => {
        console.log(message.body);
    }
);

subscriptions.notifications = client.subscribe(
    '/topic/notifications',
    (message) => {
        console.log(message.body);
    }
);

Массовая отписка

Object.values(subscriptions).forEach(subscription => {
    subscription.unsubscribe();
});

Подписка на очереди

Queue

Очереди обычно используются для распределения задач.

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

    console.log('Задача:', message.body);

});

Подписка на топики

Topic

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

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

    console.log('Новость:', message.body);

});

Пользовательские каналы

User destination

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

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

    console.log(message.body);

});

Подписки в React

useEffect и cleanup

Одна из самых распространённых ошибок — отсутствие cleanup-функции.

useEffect(() => {

    const subscription = client.subscribe(
        '/topic/chat',
        (message) => {
            console.log(message.body);
        }
    );

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

}, []);

Повторный рендер и подписки

Без cleanup React-компонент может создавать новые подписки после каждого рендера.

Это приводит к лавинообразному росту обработчиков.


Подписки в Vue

onMounted и onUnmounted

let subscription = null;

onMounted(() => {

    subscription = client.subscribe(
        '/topic/chat',
        (message) => {
            console.log(message.body);
        }
    );

});

onUnmounted(() => {

    if (subscription) {
        subscription.unsubscribe();
    }

});

Подписки в Angular

ngOnInit и ngOnDestroy

private subscription: StompSubscription;

ngOnInit(): void {

    this.subscription = this.client.subscribe(
        '/topic/chat',
        (message) => {
            console.log(message.body);
        }
    );

}

ngOnDestroy(): void {

    if (this.subscription) {
        this.subscription.unsubscribe();
    }

}

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

Автоматический reconnect

После переподключения брокер может потерять старые подписки.


Повторная регистрация

client.onConn ect = () => {

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

        console.log(message.body);

    });

};

Подписки обычно создаются именно внутри onConnect.


Временные подписки

Подписка на короткий период

const subscription = client.subscribe(
    '/topic/temp',
    (message) => {
        console.log(message.body);
    }
);

setTimeout(() => {

    subscription.unsubscribe();

}, 5000);

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

Локальная фильтрация

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

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

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

    console.log(data);

});

Подписка с уникальным id

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

client.subscribe(
    '/topic/chat',
    (message) => {
        console.log(message.body);
    },
    {
        id: 'chat-subscription'
    }
);

Для чего нужен id

Идентификатор применяется:

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

Бинарные сообщения

Работа с binaryBody

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

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

    const binary = message.binaryBody;

    console.log(binary);

});

Производительность подписок

Частые сообщения

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


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

let counter = 0;

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

    counter++;

    if (counter % 10 !== 0) {
        return;
    }

    console.log(message.body);

});

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

Накопление данных

const buffer = [];

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

    buffer.push(message.body);

});

Пакетная обработка

setInterval(() => {

    if (!buffer.length) {
        return;
    }

    console.log(buffer.splice(0));

}, 1000);

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

Опасность необработанных исключений

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

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

    saveOrder(data);

});

Любая ошибка внутри callback способна нарушить поток обработки.


Безопасная обработка

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

    try {

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

        saveOrder(data);

    } catch (error) {

        console.error(error);

    }

});

Асинхронная обработка сообщений

Async callback

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

    try {

        await processTask(message.body);

        message.ack();

    } catch (error) {

        message.nack();

    }

}, {
    ack: 'client-individual'
});

Ограничение количества подписок

Массовые подключения

Тысячи подписок могут привести к серьёзным проблемам:

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

Оптимизация

Часто выгоднее:

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