Наблюдатель для сообщений

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

В классической реализации наблюдателя существует субъект, который хранит список подписчиков и уведомляет их при изменении состояния. В STOMP.js роль субъекта выполняет брокер сообщений (RabbitMQ, ActiveMQ, Kafka через мосты), а клиентская библиотека лишь управляет локальными подписками и маршрутизацией полученных кадров (frames).

Подписка как наблюдатель

Каждый вызов подписки в STOMP.js создает объект, инкапсулирующий обработчик сообщений:

const subscription = client.subscribe('/topic/orders', (message) => {
  const payload = JSON.parse(message.body);
  console.log('Получено сообщение:', payload);
});

В этом контексте callback-функция выступает наблюдателем. Она получает уведомления каждый раз, когда брокер публикует новое сообщение в указанный канал. Подписка становится активным наблюдателем, а возвращаемый объект — управляющим дескриптором жизненного цикла наблюдения.

Ключевая особенность заключается в том, что STOMP.js не навязывает структуру управления подписчиками. Разработчик сам формирует архитектуру наблюдателей поверх низкоуровневого API.

Множественные наблюдатели и фан-аут модель

Одно из фундаментальных свойств STOMP — поддержка fan-out доставки сообщений. Несколько подписчиков могут одновременно наблюдать один и тот же канал:

const sub1 = client.subscribe('/topic/chat', handleUserMessages);
const sub2 = client.subscribe('/topic/chat', handleAnalytics);
const sub3 = client.subscribe('/topic/chat', logMessages);

Каждый обработчик получает копию сообщения независимо от остальных. Это поведение соответствует классической модели наблюдателя, где один субъект уведомляет множество независимых слушателей.

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

Инкапсуляция наблюдателей через слой абстракции

В реальных приложениях прямое использование subscribe приводит к рассеиванию логики обработки сообщений. Для построения устойчивой архитектуры вводится слой наблюдателей, который централизует управление подписками.

Пример базового менеджера наблюдателей:

class MessageObserverManager {
  constructor(stompClient) {
    this.client = stompClient;
    this.subscriptions = new Map();
  }

  addObserver(topic, handler) {
    if (!this.subscriptions.has(topic)) {
      const subscription = this.client.subscribe(topic, (message) => {
        const data = JSON.parse(message.body);
        this.notify(topic, data);
      });

      this.subscriptions.set(topic, {
        subscription,
        observers: new Set()
      });
    }

    this.subscriptions.get(topic).observers.add(handler);
  }

  notify(topic, data) {
    const entry = this.subscriptions.get(topic);
    if (!entry) return;

    entry.observers.forEach(handler => handler(data));
  }

  removeObserver(topic, handler) {
    const entry = this.subscriptions.get(topic);
    if (!entry) return;

    entry.observers.delete(handler);

    if (entry.observers.size === 0) {
      entry.subscription.unsubscribe();
      this.subscriptions.delete(topic);
    }
  }
}

Здесь STOMP-подписка становится низкоуровневым транспортом, а наблюдатели управляются локально. Это позволяет реализовать паттерн Observer независимо от ограничений библиотеки.

Фильтрация событий внутри наблюдателей

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

manager.addObserver('/topic/orders', (data) => {
  if (data.status === 'PAID') {
    console.log('Оплаченный заказ:', data.id);
  }
});

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

Диспетчеризация наблюдателей

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

class EventDispatcher {
  constructor() {
    this.channels = new Map();
  }

  register(eventType, handler) {
    if (!this.channels.has(eventType)) {
      this.channels.set(eventType, new Set());
    }

    this.channels.get(eventType).add(handler);
  }

  dispatch(eventType, payload) {
    const handlers = this.channels.get(eventType);
    if (!handlers) return;

    handlers.forEach(fn => fn(payload));
  }
}

Интеграция с STOMP.js выполняется через единый входной поток:

client.subscribe('/topic/orders', (message) => {
  const event = JSON.parse(message.body);
  dispatcher.dispatch(event.type, event);
});

Такой подход отделяет транспортный уровень от логики наблюдателей и формирует событийную архитектуру.

Управление жизненным циклом наблюдателей

Ключевая проблема наблюдателей в STOMP.js — утечки подписок при длительной работе приложения. Каждый subscribe создает активное соединение, которое должно быть корректно завершено.

Жизненный цикл наблюдателя включает три состояния:

  • регистрация обработчика
  • активное получение сообщений
  • удаление и освобождение ресурсов

Корректная деактивация:

const handler = (data) => {
  console.log(data);
};

manager.addObserver('/topic/notifications', handler);

// позже
manager.removeObserver('/topic/notifications', handler);

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

Наблюдатели и конкурентность обработки

STOMP.js работает асинхронно, и каждый наблюдатель выполняется независимо. Это создает потенциальные состояния гонки при изменении общего состояния приложения.

Пример проблемного сценария:

client.subscribe('/topic/state', (message) => {
  globalState.value = JSON.parse(message.body);
});

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

Композиция наблюдателей

Сложные системы используют композицию наблюдателей, когда один обработчик делегирует работу нескольким специализированным функциям:

function createOrderObservers() {
  return {
    onCreated: (data) => console.log('Создан заказ', data),
    onUpdated: (data) => console.log('Обновление заказа', data),
    onDeleted: (data) => console.log('Удален заказ', data)
  };
}

Далее диспетчер маршрутизирует события:

const observers = createOrderObservers();

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

  if (event.type === 'created') observers.onCreated(event);
  if (event.type === 'updated') observers.onUpdated(event);
  if (event.type === 'deleted') observers.onDeleted(event);
});

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

Иерархия наблюдателей в клиентской архитектуре

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

  • транспортный уровень (STOMP subscribe)
  • уровень маршрутизации событий
  • бизнес-обработчики
  • UI-реакции

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

Повторное использование наблюдателей

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

function createLogger(prefix) {
  return (data) => {
    console.log(`[${prefix}]`, data);
  };
}

client.subscribe('/topic/a', createLogger('A'));
client.subscribe('/topic/b', createLogger('B'));

Такой подход превращает наблюдателей в фабрики функций, повышая гибкость и снижая дублирование кода.

Отложенные и условные наблюдатели

Некоторые сценарии требуют активации наблюдателей только при определенных условиях:

let active = false;

client.subscribe('/topic/metrics', (message) => {
  if (!active) return;

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

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

Событийная консистентность наблюдателей

При построении системы наблюдателей важно учитывать, что STOMP не гарантирует доставку сообщений в строгом порядке при распределенных брокерах. Это влияет на консистентность состояния наблюдателей.

Решение заключается в использовании версионирования сообщений:

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

  if (event.version < lastVersion) return;
  lastVersion = event.version;

  handleEvent(event);
});

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