Множественные подписки

В STOMP.js клиент может одновременно подписываться на большое количество каналов, очередей и пользовательских направлений. Такой подход используется в системах реального времени, где приложение получает данные из разных источников:

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

Множественные подписки позволяют одному WebSocket-соединению обслуживать сразу несколько потоков сообщений без необходимости создавать отдельный сокет для каждого канала.


Базовая схема нескольких подписок

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

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

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

client.onConn ect = () => {

    client.subscribe('/topic/news', (message) => {
        console.log('Новости:', message.body);
    });

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

    client.subscribe('/queue/tasks', (message) => {
        console.log('Задачи:', message.body);
    });
};

client.activate();

После подключения STOMP.js создаёт несколько независимых подписок внутри одного TCP/WebSocket-соединения.


Независимость подписок

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

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

Подписки не мешают друг другу.

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

client.subscribe('/topic/payments', onPaymentMessage);

client.subscribe('/topic/errors', onErrorMessage);

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


Объект Subscription

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

const subscription = client.subscribe(
    '/topic/events',
    callback
);

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

{
    id: 'sub-0',
    unsubscribe: Function
}

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


Хранение нескольких подписок

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

Хранение в объекте

const subscriptions = {};

subscriptions.news = client.subscribe(
    '/topic/news',
    onNews
);

subscriptions.chat = client.subscribe(
    '/topic/chat',
    onChat
);

subscriptions.tasks = client.subscribe(
    '/queue/tasks',
    onTasks
);

Хранение в Map

const subscriptions = new Map();

subscriptions.set(
    'metrics',
    client.subscribe('/topic/metrics', onMetrics)
);

subscriptions.set(
    'alerts',
    client.subscribe('/topic/alerts', onAlerts)
);

Map особенно удобен при динамическом управлении большим числом подписок.


Отписка от отдельных каналов

Каждая подписка удаляется независимо.

const chatSubscription = client.subscribe(
    '/topic/chat',
    onChatMessage
);

chatSubscription.unsubscribe();

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


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

Иногда требуется удалить сразу все подписки.

Через массив

const subscriptions = [];

subscriptions.push(
    client.subscribe('/topic/a', callbackA)
);

subscriptions.push(
    client.subscribe('/topic/b', callbackB)
);

subscriptions.push(
    client.subscribe('/topic/c', callbackC)
);

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

Через объект

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

Динамическое создание подписок

Во многих приложениях каналы заранее неизвестны.

Например:

  • пользователь открывает новые комнаты чата;
  • оператор подключает новые очереди;
  • система автоматически создаёт персональные каналы.

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

function subscribeToRoom(roomId) {

    return client.subscribe(
        `/topic/rooms/${roomId}`,
        (message) => {

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

            console.log(roomId, data);
        }
    );
}

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

const roomSubscriptions = {};

const roomIds = [101, 102, 103];

roomIds.forEach((roomId) => {

    roomSubscriptions[roomId] = client.subscribe(
        `/topic/rooms/${roomId}`,
        (message) => {

            console.log(
                `Комната ${roomId}:`,
                message.body
            );
        }
    );
});

Динамическая отписка

function unsubscribeRoom(roomId) {

    const subscription = roomSubscriptions[roomId];

    if (!subscription) {
        return;
    }

    subscription.unsubscribe();

    delete roomSubscriptions[roomId];
}

Автоматическое создание callback-функций

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

function createHandler(channel) {

    return (message) => {

        console.log(
            `[${channel}]`,
            message.body
        );
    };
}

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

client.subscribe(
    '/topic/payments',
    createHandler('payments')
);

Разделение подписок по модулям

В больших приложениях подписки группируются логически.

Пример структуры

const chatModule = {
    subscriptions: []
};

const notificationModule = {
    subscriptions: []
};

const analyticsModule = {
    subscriptions: []
};

Инициализация подписок модуля

chatModule.subscriptions.push(
    client.subscribe('/topic/chat/global', onGlobalChat)
);

chatModule.subscriptions.push(
    client.subscribe('/topic/chat/private', onPrivateChat)
);

Очистка подписок модуля

function destroyModule(module) {

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

    module.subscriptions = [];
}

Множественные подписки и ACK

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

client.subscribe(
    '/queue/orders',
    onOrders,
    { ack: 'client' }
);

client.subscribe(
    '/topic/logs',
    onLogs,
    { ack: 'auto' }
);

Это позволяет:

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

Изолированная обработка ошибок

Ошибка в callback не должна нарушать работу остальных подписок.

Неправильный подход

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

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

    processData(data);
});

Если JSON повреждён — callback завершится исключением.


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

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

    try {

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

        processData(data);

    } catch (error) {

        console.error(error);
    }
});

Производительность при большом числе подписок

Количество подписок влияет на:

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

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

  • утечки подписок;
  • повторные подписки;
  • дублирование callback;
  • количество активных destinations.

Проблема дублирующихся подписок

Частая ошибка — повторная подписка на один и тот же канал.

Ошибочный код

function init() {

    client.subscribe('/topic/news', onNews);
}

Если init() вызывается многократно, создаются новые подписки.


Правильный вариант

let newsSubscription = null;

function init() {

    if (newsSubscription) {
        return;
    }

    newsSubscription = client.subscribe(
        '/topic/news',
        onNews
    );
}

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

После переподключения STOMP.js автоматически восстанавливает подписки, созданные через активный клиент.

Однако существуют особенности:

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

Пересоздание подписок вручную

Иногда приложение полностью контролирует процесс reconnect.

const channels = [
    '/topic/a',
    '/topic/b',
    '/topic/c'
];

client.onConn ect = () => {

    channels.forEach((channel) => {

        client.subscribe(channel, (message) => {

            console.log(channel, message.body);
        });
    });
};

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


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

Во многих системах используются персональные направления.

client.subscribe(
    '/user/queue/notifications',
    onNotification
);

Одновременно могут существовать:

client.subscribe('/user/queue/messages', onMessages);

client.subscribe('/user/queue/tasks', onTasks);

client.subscribe('/user/queue/events', onEvents);

Комбинирование queue и topic

Одно соединение может работать с разными типами destinations.

client.subscribe('/topic/system', onSystem);

client.subscribe('/queue/jobs', onJobs);

client.subscribe('/exchange/logs', onLogs);

Это особенно распространено в RabbitMQ STOMP Adapter.


Идентификаторы подписок

STOMP требует уникальный id для каждой подписки.

STOMP.js генерирует идентификаторы автоматически:

sub-0
sub-1
sub-2

Но можно задавать собственные значения.

client.subscribe(
    '/topic/news',
    onNews,
    {
        id: 'news-subscription'
    }
);

Ручное управление ID подписок

Ручные идентификаторы полезны:

  • при отладке;
  • в логировании;
  • при диагностике reconnect;
  • в сложных распределённых системах.
client.subscribe(
    '/topic/orders',
    onOrders,
    {
        id: 'orders-v1'
    }
);

Коллизии идентификаторов

Нельзя использовать одинаковый id для нескольких подписок.

Ошибка

client.subscribe(
    '/topic/a',
    callbackA,
    { id: 'shared-id' }
);

client.subscribe(
    '/topic/b',
    callbackB,
    { id: 'shared-id' }
);

Результат зависит от брокера:

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

Управление подписками в SPA

В React, Vue и Angular особенно важно удалять подписки при уничтожении компонентов.

Проблема утечки

function openPage() {

    client.subscribe('/topic/page', callback);
}

При повторном открытии страницы количество подписок растёт.


Контролируемая очистка

let subscription = null;

function mount() {

    subscription = client.subscribe(
        '/topic/page',
        callback
    );
}

function unmount() {

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

Централизованный менеджер подписок

В крупных проектах создаётся единый слой управления.

Пример менеджера

class SubscriptionManager {

    constructor(client) {

        this.client = client;

        this.items = new Map();
    }

    subscribe(name, destination, callback) {

        if (this.items.has(name)) {
            return;
        }

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

        this.items.set(name, subscription);
    }

    unsubscribe(name) {

        const subscription = this.items.get(name);

        if (!subscription) {
            return;
        }

        subscription.unsubscribe();

        this.items.delete(name);
    }

    unsubscribeAll() {

        this.items.forEach((subscription) => {
            subscription.unsubscribe();
        });

        this.items.clear();
    }
}

Использование менеджера

const manager = new SubscriptionManager(client);

manager.subscribe(
    'orders',
    '/topic/orders',
    onOrders
);

manager.subscribe(
    'payments',
    '/topic/payments',
    onPayments
);

manager.unsubscribe('orders');

Логирование множественных подписок

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

function subscribeWithLog(destination, callback) {

    console.log(
        'SUBSCRIBE:',
        destination
    );

    return client.subscribe(destination, callback);
}

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

console.log(
    'Количество подписок:',
    subscriptions.size
);

Это помогает находить:

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

Практическая архитектура множественных подписок

Распространённая схема:

const subscriptions = {
    system: [],
    user: [],
    chat: [],
    analytics: []
};

Разделение по категориям упрощает:

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

Массовая подписка через конфигурацию

const configs = [
    {
        key: 'orders',
        destination: '/topic/orders',
        handler: onOrders
    },
    {
        key: 'payments',
        destination: '/topic/payments',
        handler: onPayments
    },
    {
        key: 'alerts',
        destination: '/topic/alerts',
        handler: onAlerts
    }
];

configs.forEach((config) => {

    subscriptions[config.key] = client.subscribe(
        config.destination,
        config.handler
    );
});

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

Брокеры сообщений могут ограничивать:

  • число подписок на соединение;
  • скорость создания подписок;
  • число пользовательских очередей.

При высоких нагрузках используются:

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

Агрегация сообщений

Вместо множества подписок:

/topic/orders/created
/topic/orders/updated
/topic/orders/deleted

используется один канал:

/topic/orders

Тип события передаётся внутри сообщения:

{
    "type": "updated",
    "payload": {}
}

Такой подход уменьшает количество подписок.


Баланс между числом подписок и размером трафика

Два крайних подхода:

Много специализированных каналов

Плюсы:

  • меньше лишнего трафика;
  • проще фильтрация.

Минусы:

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

Один агрегированный канал

Плюсы:

  • меньше подписок;
  • проще reconnect.

Минусы:

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

Архитектура выбирается исходя из:

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