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

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

Без серверной фильтрации брокер вынужден передавать все сообщения всем подписчикам канала, а фильтрация переносится на клиентскую сторону. Такой подход создаёт ряд проблем:

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

При фильтрации на уровне брокера сообщение получает только тот клиент, который соответствует заданным условиям.


Клиентская фильтрация и её недостатки

Простейший вариант выглядит следующим образом:

client.subscribe('/topic/orders', (message) => {
    const order = JSON.parse(message.body);

    if (order.userId !== currentUserId) {
        return;
    }

    renderOrder(order);
});

В этом примере:

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

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

Особенно заметны проблемы в системах:

  • онлайн-чатов;
  • биржевых терминалов;
  • игровых серверов;
  • систем мониторинга;
  • IoT-платформ;
  • систем уведомлений.

Серверная фильтрация через отдельные destination

Наиболее распространённый способ фильтрации — разделение сообщений по destination.

Например:

/topic/orders/user/15
/topic/orders/user/28
/topic/orders/user/41

Подписка выполняется только на собственный канал:

client.subscribe('/topic/orders/user/15', (message) => {
    const order = JSON.parse(message.body);

    renderOrder(order);
});

Теперь брокер:

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

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

client.publish({
    destination: '/topic/orders/user/15',
    body: JSON.stringify({
        id: 1001,
        status: 'created'
    })
});

Иерархическая маршрутизация

Многие брокеры поддерживают иерархическую структуру адресов.

Пример:

/topic/orders/europe/kz
/topic/orders/europe/de
/topic/orders/asia/jp

Подписчики получают только нужную географическую категорию:

client.subscribe('/topic/orders/europe/kz', callback);

Такой подход позволяет:

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

Фильтрация через пользовательские очереди

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

Пример для RabbitMQ:

/exchange/amq.direct/user.15

Подписка:

client.subscribe('/exchange/amq.direct/user.15', (message) => {
    console.log(message.body);
});

Отправка:

client.publish({
    destination: '/exchange/amq.direct/user.15',
    body: 'Private message'
});

Такой механизм часто применяется:

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

Topic-based filtering

Фильтрация может выполняться через структуру topic.

Например:

/topic/news/sport
/topic/news/politics
/topic/news/finance

Клиент подписывается только на нужную категорию:

client.subscribe('/topic/news/finance', (message) => {
    console.log('Finance:', message.body);
});

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

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

Wildcard-подписки

Некоторые брокеры поддерживают wildcard-маршрутизацию.

Примеры:

/topic/orders/*
/topic/orders/**

Либо:

topic.orders.*
topic.orders.#

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

  • RabbitMQ;
  • ActiveMQ;
  • Apollo;
  • Artemis;
  • Spring Broker Relay.

Пример подписки:

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

Wildcard-маршруты позволяют:

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

Header-based filtering

Некоторые брокеры поддерживают фильтрацию по заголовкам сообщений.

Сообщение:

client.publish({
    destination: '/topic/orders',
    headers: {
        region: 'kz',
        priority: 'high'
    },
    body: JSON.stringify({
        id: 500
    })
});

Подписчик может быть настроен брокером так, чтобы получать только сообщения:

region = kz
priority = high

Подобная фильтрация особенно популярна в:

  • ActiveMQ;
  • JMS;
  • Artemis;
  • enterprise-системах.

Selector-based filtering

В JMS-совместимых брокерах используются selectors.

Пример selector:

region = 'kz' AND priority = 'high'

STOMP-подписка:

client.subscribe(
    '/topic/orders',
    callback,
    {
        selector: "region = 'kz'"
    }
);

Либо:

client.subscribe(
    '/topic/orders',
    callback,
    {
        selector: "priority = 'critical'"
    }
);

Брокер анализирует заголовки сообщений и самостоятельно решает:

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

Пример фильтрации по типу события

Сообщения:

client.publish({
    destination: '/topic/events',
    headers: {
        type: 'payment'
    },
    body: JSON.stringify({
        amount: 100
    })
});

Подписка:

client.subscribe(
    '/topic/events',
    (message) => {
        console.log(message.body);
    },
    {
        selector: "type = 'payment'"
    }
);

Сообщения другого типа не будут доставлены клиенту.


Message selectors в ActiveMQ

ActiveMQ поддерживает сложные селекторы.

Пример:

client.subscribe(
    '/queue/tasks',
    callback,
    {
        selector: "priority > 5 AND region = 'kz'"
    }
);

Допустимы:

  • логические операции;
  • числовые сравнения;
  • LIKE;
  • BETWEEN;
  • IN;
  • IS NULL.

Пример:

{
    selector: "department IN ('sales', 'support')"
}

Фильтрация по приоритету

Сообщение:

client.publish({
    destination: '/topic/alerts',
    headers: {
        priority: 'critical'
    },
    body: 'Disk failure'
});

Подписка:

client.subscribe(
    '/topic/alerts',
    callback,
    {
        selector: "priority = 'critical'"
    }
);

Это особенно полезно в:

  • системах мониторинга;
  • alerting-платформах;
  • системах безопасности;
  • DevOps-инфраструктуре.

Фильтрация по tenant

В multi-tenant системах брокер часто фильтрует сообщения по идентификатору клиента.

Публикация:

client.publish({
    destination: '/topic/data',
    headers: {
        tenant: 'companyA'
    },
    body: JSON.stringify({
        report: 'monthly'
    })
});

Подписка:

client.subscribe(
    '/topic/data',
    callback,
    {
        selector: "tenant = 'companyA'"
    }
);

Это предотвращает:

  • утечку данных;
  • смешивание tenant-трафика;
  • получение чужих событий.

Content-based routing

Некоторые брокеры поддерживают content-based routing.

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

{
    "country": "kz",
    "priority": "high",
    "amount": 5000
}

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

В чистом STOMP подобный механизм встречается редко, однако часто реализуется:

  • через промежуточные сервисы;
  • custom broker plugins;
  • интеграционные шины.

Broker relay и фильтрация

При использовании Spring WebSocket Broker Relay STOMP.js взаимодействует не со встроенным брокером Spring, а с полноценным сервером сообщений:

  • RabbitMQ;
  • ActiveMQ;
  • Artemis.

В этом случае становятся доступны:

  • selectors;
  • wildcard routing;
  • exchange-маршрутизация;
  • topic-фильтрация;
  • durable subscriptions.

Фильтрация в RabbitMQ

RabbitMQ использует routing key и exchange.

Пример:

exchange: orders
routing key: kz.high

Подписка:

client.subscribe('/exchange/orders/kz.high', callback);

Публикация:

client.publish({
    destination: '/exchange/orders/kz.high',
    body: JSON.stringify({
        id: 10
    })
});

Topic exchange в RabbitMQ

Topic exchange поддерживает шаблоны.

Примеры routing key:

kz.payment
kz.alert
us.payment

Подписка:

client.subscribe('/exchange/orders/kz.*', callback);

Либо:

client.subscribe('/exchange/orders/#', callback);

Durable subscriptions и фильтрация

Durable subscription сохраняет подписку даже после отключения клиента.

Пример:

client.subscribe(
    '/topic/orders',
    callback,
    {
        id: 'orders-subscription',
        durable: 'true',
        auto-delete: 'false'
    }
);

Вместе с selectors это позволяет:

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

Фильтрация через виртуальные хосты

Некоторые брокеры разделяют трафик через virtual hosts.

Примеры:

/dev/orders
/test/orders
/prod/orders

Либо разные vhost в RabbitMQ:

vhost-dev
vhost-prod

Такое разделение помогает:

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

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

Фильтрация на уровне брокера значительно уменьшает:

  • объём сетевого трафика;
  • количество обработчиков;
  • нагрузку на CPU браузера;
  • потребление памяти;
  • число операций JSON.parse.

Особенно это важно при:

  • десятках тысяч сообщений;
  • WebSocket-шлюзах;
  • real-time аналитике;
  • потоковой телеметрии.

Ошибки проектирования

Передача всех событий всем клиентам

Плохой пример:

/topic/all-events

Все клиенты получают:

  • сообщения чатов;
  • уведомления;
  • системные события;
  • данные мониторинга.

Это быстро приводит к перегрузке.


Фильтрация только в браузере

Плохой подход:

if (message.userId !== currentUserId) {
    return;
}

Проблемы:

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

Слишком глубокая иерархия routing key

Плохо:

/orders/europe/kz/almaty/shop/15/device/mobile

Избыточная детализация:

  • усложняет маршрутизацию;
  • замедляет broker matching;
  • затрудняет поддержку.

Безопасность фильтрации

Фильтрация не должна заменять авторизацию.

Даже если используется отдельный topic:

/topic/private/user15

Брокер обязан проверять:

  • права доступа;
  • ACL;
  • authentication;
  • разрешения на subscribe/send.

Иначе пользователь сможет подписаться на чужой канал.


ACL и ограничения подписок

Многие брокеры поддерживают ACL.

Пример концепции:

user15 -> subscribe -> /topic/private/user15
DENY -> /topic/private/*

Такой механизм критически важен для:

  • финансовых систем;
  • корпоративных платформ;
  • медицинских сервисов;
  • административных панелей.

Масштабирование через фильтрацию

Грамотно спроектированная фильтрация позволяет:

  • обслуживать больше подключений;
  • уменьшать нагрузку на WebSocket-сервер;
  • распределять брокеры по shard;
  • строить event-driven архитектуру;
  • снижать стоимость инфраструктуры.

В крупных системах фильтрация сообщений на уровне брокера является обязательным элементом архитектуры, а не дополнительной оптимизацией.