Angular-сервисы представляют собой основной механизм инкапсуляции бизнес-логики и взаимодействия с внешними источниками данных, включая WebSocket-соединения. При работе с STOMP.js в Angular сервис становится центральным слоем, управляющим жизненным циклом соединения, подписками на каналы и распределением сообщений между компонентами через реактивные потоки RxJS.
STOMP-протокол поверх WebSocket используется для обмена сообщениями в режиме publish/subscribe. STOMP.js реализует клиентскую часть протокола, обеспечивая подключение к брокерам сообщений (RabbitMQ, ActiveMQ, Spring WebSocket и другие). В Angular это обычно оформляется в виде singleton-сервиса, предоставляемого на уровне root, что гарантирует единственное соединение на всё приложение.
Сервис выступает посредником между WebSocket-слоем и UI-компонентами. Его задачи включают:
Такой подход исключает прямую работу компонентов с низкоуровневым API STOMP.js, снижая связность и упрощая тестирование.
Инкапсуляция STOMP-клиента обычно строится вокруг класса сервиса, использующего BehaviorSubject для хранения состояния соединения и потоков сообщений.
import { Injectable } from '@angular/core';
import { BehaviorSubject, Observable } from 'rxjs';
import { Client, IMessage, StompSubscription } from '@stomp/stompjs';
@Injectable({ providedIn: 'root' })
export class StompService {
private client: Client | null = null;
private connectionState$ = new BehaviorSubject<boolean>(false);
get connectionStatus(): Observable<boolean> {
return this.connectionState$.asObservable();
}
}
Состояние соединения отделяется от логики доставки сообщений, что позволяет UI реагировать на подключение независимо от бизнес-данных.
STOMP.js использует объект Client, который настраивается перед активацией соединения.
import { Client } from '@stomp/stompjs';
import * as SockJS from 'sockjs-client';
private createClient(): Client {
return new Client({
webSocketFactory: () => new SockJS('https://example.com/ws'),
reconnectDelay: 5000,
heartbeatIncoming: 4000,
heartbeatOutgoing: 4000,
debug: () => {}
});
}
Ключевые параметры:
webSocketFactory — обеспечивает совместимость с
SockJSreconnectDelay — автоматическое восстановление
соединенияheartbeatIncoming / heartbeatOutgoing — контроль
“живости” каналаdebug — отладочный выводСервис управляет активацией и деактивацией клиента, синхронизируя состояние с RxJS.
connect(): void {
if (this.client?.active) {
return;
}
this.client = this.createClient();
this.client.onConn ect = () => {
this.connectionState$.next(true);
};
this.client.onDisconn ect = () => {
this.connectionState$.next(false);
};
this.client.onStompEr ror = () => {
this.connectionState$.next(false);
};
this.client.activate();
}
Данная схема обеспечивает централизованную реакцию на события транспортного уровня.
STOMP использует модель destinations (например,
/topic/chat, /queue/events). Каждая подписка
возвращает объект StompSubscription, который необходимо хранить для
последующей отписки.
private subscriptions: Map<string, StompSubscription> = new Map();
subscribe<T>(destination: string): Observable<T> {
return new Observable<T>(observer => {
if (!this.client?.connected) {
observer.error('STOMP not connected');
return;
}
const sub = this.client.subscribe(destination, (message: IMessage) => {
const body = JSON.parse(message.body);
observer.next(body as T);
});
this.subscriptions.set(destination, sub);
return () => {
sub.unsubscribe();
this.subscriptions.delete(destination);
};
});
}
Такой подход связывает STOMP-подписку с жизненным циклом Observable, автоматически освобождая ресурсы при отписке компонентов.
Отправка данных в брокер осуществляется через метод publish или send (в зависимости от версии STOMP.js).
send(destination: string, payload: unknown): void {
if (!this.client?.connected) {
return;
}
this.client.publish({
destination,
body: JSON.stringify(payload)
});
}
При необходимости можно добавлять заголовки:
this.client.publish({
destination,
body: JSON.stringify(payload),
headers: {
priority: 'high',
authorization: 'Bearer token'
}
});
Для распространения сообщений между компонентами используется Subject или BehaviorSubject.
private messages$ = new BehaviorSubject<any[]>([]);
listenToTopic(): void {
this.client?.subscribe('/topic/updates', (msg) => {
const data = JSON.parse(msg.body);
const current = this.messages$.value;
this.messages$.next([...current, data]);
});
}
Такой подход позволяет реализовать глобальное состояние без дополнительного state-management слоя.
В сложных приложениях используется разделение потоков:
private userMessages$ = new Subject<any>();
private systemEvents$ = new Subject<any>();
private errorEvents$ = new Subject<any>();
Роутинг сообщений выполняется на уровне сервиса:
private routeMessage(destination: string, payload: any): void {
if (destination.includes('/topic/system')) {
this.systemEvents$.next(payload);
} else if (destination.includes('/topic/user')) {
this.userMessages$.next(payload);
}
}
STOMP.js поддерживает автоматическое переподключение, однако состояние подписок требует восстановления вручную.
private restoreSubscriptions(): void {
this.subscriptions.forEach((_, destination) => {
this.subscribe(destination).subscribe();
});
}
Вызывается после успешного reconnection:
this.client.onConn ect = () => {
this.connectionState$.next(true);
this.restoreSubscriptions();
};
WebSocket-события происходят вне Angular Zone, что может приводить к отсутствию обновления UI.
constructor(private ngZone: NgZone) {}
this.client.subscribe('/topic/data', (msg) => {
this.ngZone.run(() => {
const data = JSON.parse(msg.body);
this.messages$.next(data);
});
});
Это гарантирует корректный запуск change detection.
При использовании защищённых брокеров сообщений токен передаётся через headers CONNECT:
this.client = new Client({
webSocketFactory: () => new SockJS('/ws'),
connectHeaders: {
Authorization: `Bearer ${token}`
}
});
При обновлении токена требуется пересоздание соединения, так как STOMP не всегда поддерживает динамическое обновление headers.
Для повышения надёжности используется строгая типизация payload:
interface ChatMessage {
id: string;
text: string;
author: string;
timestamp: number;
}
subscribe<ChatMessage>('/topic/chat');
Это снижает количество runtime-ошибок при обработке JSON.
При тестировании STOMP-сервиса используется мокирование Client:
const mockClient = {
activate: jasmine.createSpy(),
deactivate: jasmine.createSpy(),
subscribe: jasmine.createSpy().and.returnValue({
unsubscribe: jasmine.createSpy()
})
};
Сервис проверяется как изолированная единица без реального WebSocket.
Критическим аспектом является корректное освобождение ресурсов:
disconnect(): void {
this.subscriptions.forEach(sub => sub.unsubscribe());
this.subscriptions.clear();
this.client?.deactivate();
this.connectionState$.next(false);
}
При росте приложения STOMP-сервис часто расширяется до нескольких уровней:
Такое разделение позволяет:
В некоторых системах требуется несколько брокеров или пространств:
private clients: Map<string, Client> = new Map();
Каждый клиент управляется отдельно, а маршрутизация сообщений выполняется по ключу соединения.
При высокой частоте сообщений используется буферизация:
private buffered$ = new Subject<any>();
this.client.subscribe('/topic/high-load', msg => {
this.buffered$.next(JSON.parse(msg.body));
});
Далее применяются RxJS операторы bufferTime,
throttleTime, auditTime для контроля нагрузки
UI.
STOMP-сервис в Angular формирует слой:
Он объединяет асинхронный поток сообщений с реактивной моделью Angular, обеспечивая предсказуемую обработку событий, управляемые подписки и централизованный контроль соединения.