Дублирование сообщений — одна из наиболее распространённых проблем при работе с WebSocket-соединениями и брокерами сообщений через STOMP.js. Ошибка может проявляться по-разному:
Подобные ситуации особенно критичны для:
Наиболее частая причина дублирования — повторное создание подписки после восстановления соединения.
import { Client } fr om '@stomp/stompjs';
const client = new Client({
brokerURL: 'ws://localhost:15674/ws',
reconnectDelay: 5000
});
function subscribeToChat() {
client.subscribe('/topic/chat', (message) => {
console.log('Новое сообщение:', message.body);
});
}
client.onConn ect = () => {
subscribeToChat();
};
client.activate();
На первый взгляд код выглядит корректным. Проблема возникает при реконнекте.
Если соединение разорвётся:
onConnect.subscribeToChat().В результате сервер будет отправлять одно сообщение сразу нескольким подписчикам.
Каждый вызов:
client.subscribe(...)
создаёт новую подписку.
STOMP.js не заменяет предыдущую автоматически.
Например:
client.subscribe('/topic/chat', callback);
client.subscribe('/topic/chat', callback);
client.subscribe('/topic/chat', callback);
создаёт три независимые подписки.
При публикации:
client.publish({
destination: '/topic/chat',
body: 'Hello'
});
callback выполнится три раза.
Каждая подписка возвращает объект Subscription.
let chatSubscription = null;
client.onConn ect = () => {
if (chatSubscription) {
chatSubscription.unsubscribe();
}
chatSubscription = client.subscribe('/topic/chat', (message) => {
console.log(message.body);
});
};
Теперь старая подписка удаляется перед созданием новой.
В больших приложениях ручное управление подписками быстро становится неудобным.
class SubscriptionManager {
constructor(client) {
this.client = client;
this.subscriptions = new Map();
}
subscribe(key, destination, callback) {
if (this.subscriptions.has(key)) {
this.subscriptions.get(key).unsubscribe();
}
const subscription = this.client.subscribe(
destination,
callback
);
this.subscriptions.set(key, subscription);
return subscription;
}
unsubscribe(key) {
if (!this.subscriptions.has(key)) {
return;
}
this.subscriptions.get(key).unsubscribe();
this.subscriptions.delete(key);
}
unsubscribeAll() {
for (const subscription of this.subscriptions.values()) {
subscription.unsubscribe();
}
this.subscriptions.clear();
}
}
Использование:
const manager = new SubscriptionManager(client);
client.onConn ect = () => {
manager.subscribe(
'chat',
'/topic/chat',
(message) => {
console.log(message.body);
}
);
};
Иногда проблема связана не с подписками, а с созданием нескольких STOMP-клиентов.
function connect() {
const client = new Client({
brokerURL: 'ws://localhost:15674/ws'
});
client.onConn ect = () => {
client.subscribe('/topic/chat', (message) => {
console.log(message.body);
});
};
client.activate();
}
connect();
connect();
connect();
Будет создано:
Каждое сообщение придёт трижды.
Для предотвращения подобных ошибок обычно создают единый экземпляр клиента.
class SocketService {
static instance = null;
constructor() {
if (SocketService.instance) {
return SocketService.instance;
}
this.client = new Client({
brokerURL: 'ws://localhost:15674/ws'
});
SocketService.instance = this;
}
getClient() {
return this.client;
}
}
const socketService = new SocketService();
const client = socketService.getClient();
В React проблема особенно распространена.
useEffect(() => {
client.subscribe('/topic/chat', (message) => {
console.log(message.body);
});
}, []);
Если компонент размонтируется и смонтируется заново:
useEffect(() => {
const subscription = client.subscribe(
'/topic/chat',
(message) => {
console.log(message.body);
}
);
return () => {
subscription.unsubscribe();
};
}, []);
mounted() {
this.subscription = client.subscribe(
'/topic/chat',
this.onMessage
);
}
Если компонент пересоздаётся:
beforeUnmount() {
if (this.subscription) {
this.subscription.unsubscribe();
}
}
Режим подтверждений также способен вызывать повторную доставку сообщений.
client.subscribe(
'/queue/tasks',
onMessage,
{ ack: 'auto' }
);
Сообщение считается обработанным сразу после отправки клиенту.
client.subscribe(
'/queue/tasks',
(message) => {
processTask(message);
message.ack();
},
{ ack: 'client' }
);
Если ack() не будет вызван:
Типичный сценарий:
Такое поведение считается нормальным.
Для защиты от повторов обработка должна быть идемпотентной.
const processedMessages = new Set();
client.subscribe('/queue/tasks', (message) => {
const id = message.headers['message-id'];
if (processedMessages.has(id)) {
return;
}
processedMessages.add(id);
processTask(message.body);
message.ack();
}, {
ack: 'client'
});
Предыдущий подход опасен утечкой памяти.
Лучше использовать ограниченный кеш.
class MessageCache {
constructor(lim it = 1000) {
this.limit = limit;
this.messages = new Map();
}
has(id) {
return this.messages.has(id);
}
add(id) {
this.messages.set(id, Date.now());
if (this.messages.size > this.limit) {
const firstKey =
this.messages.keys().next().value;
this.messages.delete(firstKey);
}
}
}
Транзакции помогают уменьшить риск повторной обработки.
const tx = client.begin();
try {
processPayment();
message.ack({
transaction: tx.id
});
tx.commit();
} catch (error) {
tx.abort();
}
Иногда проблема находится не на клиенте.
Например:
При диагностике необходимо анализировать:
Heartbeat способен провоцировать ложные reconnect.
const client = new Client({
brokerURL: 'ws://localhost:15674/ws',
heartbeatIncoming: 4000,
heartbeatOutgoing: 4000,
reconnectDelay: 5000
});
Если heartbeat слишком агрессивен:
Обычно используются более безопасные интервалы:
heartbeatIncoming: 10000,
heartbeatOutgoing: 10000
или:
heartbeatIncoming: 20000,
heartbeatOutgoing: 20000
Иногда reconnect и subscribe выполняются параллельно.
client.onConn ect = () => {
subscribeChat();
setTimeout(() => {
subscribeChat();
}, 1000);
};
Будет создано две подписки.
let subscribed = false;
function subscribeChat() {
if (subscribed) {
return;
}
subscribed = true;
client.subscribe('/topic/chat', (message) => {
console.log(message.body);
});
}
STOMP поддерживает идентификаторы подписок.
client.subscribe(
'/topic/chat',
callback,
{
id: 'chat-subscription'
}
);
Некоторые брокеры предотвращают дублирование одинаковых ID.
Однако это зависит от реализации сервера.
Иногда надёжнее фильтровать сообщения самостоятельно.
Сервер:
const message = {
id: crypto.randomUUID(),
text: 'Hello'
};
Клиент:
const received = new Set();
client.subscribe('/topic/chat', (message) => {
const data = JSON.parse(message.body);
if (received.has(data.id)) {
return;
}
received.add(data.id);
renderMessage(data);
});
В распределённых системах локального Set недостаточно.
Используется Redis:
const exists = await redis.get(message.id);
if (exists) {
return;
}
await redis.set(message.id, 1, {
EX: 300
});
Большинство брокеров не гарантируют истинную модель exactly once.
На практике обычно используются:
Именно поэтому защита от дублей почти всегда реализуется на прикладном уровне.
Сообщение:
Повторов практически нет.
Но возможна потеря данных.
Сообщение гарантированно доставляется.
Но возможны:
Именно этот режим чаще всего используется совместно со STOMP.
Надёжная система обычно включает:
Очень полезно отслеживать создание подписок.
function createSubscription(destination, callback) {
console.log(
'SUBSCRIBE:',
destination,
Date.now()
);
return client.subscribe(destination, callback);
}
console.log(client);
Многие реализации STOMP.js содержат внутренние структуры:
Это помогает обнаруживать накопление подписок.
Вкладка Network позволяет увидеть:
Если в WebSocket visible multiple SUBSCRIBE frames для одного destination — проблема практически всегда связана с повторной подпиской.
const client = new Client({
brokerURL: 'ws://localhost:15674/ws',
debug(str) {
console.log(str);
}
});
STOMP.js начнёт выводить:
>>> CONNECT
>>> SUBSCRIBE
>>> SUBSCRIBE
>>> SUBSCRIBE
<<< MESSAGE
Повторяющиеся SUBSCRIBE сразу указывают на источник дублей.
import { Client } from '@stomp/stompjs';
class RealtimeService {
constructor() {
this.subscriptions = new Map();
this.client = new Client({
brokerURL: 'ws://localhost:15674/ws',
reconnectDelay: 5000
});
this.client.onConn ect = () => {
this.restoreSubscriptions();
};
}
connect() {
this.client.activate();
}
subscribe(key, destination, callback) {
if (this.subscriptions.has(key)) {
return;
}
const subscription = this.client.subscribe(
destination,
callback
);
this.subscriptions.set(key, {
destination,
callback,
subscription
});
}
restoreSubscriptions() {
for (const [key, item] of this.subscriptions) {
item.subscription.unsubscribe();
item.subscription = this.client.subscribe(
item.destination,
item.callback
);
}
}
unsubscribe(key) {
if (!this.subscriptions.has(key)) {
return;
}
this.subscriptions
.get(key)
.subscription
.unsubscribe();
this.subscriptions.delete(key);
}
}
Такой подход: