Паттерн наблюдателя естественным образом реализуется в модели работы
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.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);
});
Такой механизм позволяет наблюдателям поддерживать актуальное состояние независимо от порядка доставки событий.