Дублирование сообщений

Дублирование сообщений — одна из наиболее распространённых проблем при работе с 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();

На первый взгляд код выглядит корректным. Проблема возникает при реконнекте.

Если соединение разорвётся:

  1. STOMP.js выполнит reconnect.
  2. Снова вызовется onConnect.
  3. Снова выполнится subscribeToChat().
  4. Создастся новая подписка.

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


Как работает накопление подписок

Каждый вызов:

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 выполнится три раза.


Использование unsubscribe

Каждая подписка возвращает объект Subscription.

Правильный подход

let chatSubscription = null;

client.onConn ect = () => {

    if (chatSubscription) {
        chatSubscription.unsubscribe();
    }

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

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


Централизованный менеджер подписок

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

Пример SubscriptionManager

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();

Будет создано:

  • 3 WebSocket-соединения;
  • 3 STOMP-сессии;
  • 3 подписки.

Каждое сообщение придёт трижды.


Singleton-подключение

Для предотвращения подобных ошибок обычно создают единый экземпляр клиента.

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

В React проблема особенно распространена.

Ошибочный useEffect

useEffect(() => {

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

}, []);

Если компонент размонтируется и смонтируется заново:

  • создастся новая подписка;
  • старая останется активной.

Очистка подписки в React

useEffect(() => {

    const subscription = client.subscribe(
        '/topic/chat',
        (message) => {
            console.log(message.body);
        }
    );

    return () => {
        subscription.unsubscribe();
    };

}, []);

Дублирование в Vue

Ошибочный вариант

mounted() {

    this.subscription = client.subscribe(
        '/topic/chat',
        this.onMessage
    );
}

Если компонент пересоздаётся:

  • подписка остаётся;
  • создаётся новая.

Правильная очистка

beforeUnmount() {

    if (this.subscription) {
        this.subscription.unsubscribe();
    }
}

Дублирование при ACK

Режим подтверждений также способен вызывать повторную доставку сообщений.

AUTO ACK

client.subscribe(
    '/queue/tasks',
    onMessage,
    { ack: 'auto' }
);

Сообщение считается обработанным сразу после отправки клиенту.


CLIENT ACK

client.subscribe(
    '/queue/tasks',
    (message) => {

        processTask(message);

        message.ack();
    },
    { ack: 'client' }
);

Если ack() не будет вызван:

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

Повторная доставка после reconnect

Типичный сценарий:

  1. Клиент получает сообщение.
  2. Обработка ещё не завершена.
  3. Соединение обрывается.
  4. ACK не отправлен.
  5. После reconnect брокер повторно доставляет сообщение.

Такое поведение считается нормальным.


Idempotent-обработка сообщений

Для защиты от повторов обработка должна быть идемпотентной.

Пример фильтрации

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'
});

Ограничение памяти Set

Предыдущий подход опасен утечкой памяти.

Лучше использовать ограниченный кеш.

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);
        }
    }
}

Использование transaction

Транзакции помогают уменьшить риск повторной обработки.

const tx = client.begin();

try {

    processPayment();

    message.ack({
        transaction: tx.id
    });

    tx.commit();

} catch (error) {

    tx.abort();
}

Дублирование из-за серверной логики

Иногда проблема находится не на клиенте.

Например:

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

Проверка broker logs

При диагностике необходимо анализировать:

  • RabbitMQ logs;
  • ActiveMQ logs;
  • Artemis logs;
  • WebSocket gateway logs;
  • reverse proxy logs.

Дублирование из-за heartbeat reconnect

Heartbeat способен провоцировать ложные reconnect.

Пример

const client = new Client({

    brokerURL: 'ws://localhost:15674/ws',

    heartbeatIncoming: 4000,
    heartbeatOutgoing: 4000,

    reconnectDelay: 5000
});

Если heartbeat слишком агрессивен:

  • соединение может считаться потерянным;
  • выполняется reconnect;
  • создаются новые подписки.

Корректная настройка heartbeat

Обычно используются более безопасные интервалы:

heartbeatIncoming: 10000,
heartbeatOutgoing: 10000

или:

heartbeatIncoming: 20000,
heartbeatOutgoing: 20000

Дублирование из-за race condition

Иногда 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);
    });
}

Уникальные subscription id

STOMP поддерживает идентификаторы подписок.

client.subscribe(
    '/topic/chat',
    callback,
    {
        id: 'chat-subscription'
    }
);

Некоторые брокеры предотвращают дублирование одинаковых ID.

Однако это зависит от реализации сервера.


Дедупликация на уровне приложения

Иногда надёжнее фильтровать сообщения самостоятельно.

UUID сообщений

Сервер:

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);
});

Redis для дедупликации

В распределённых системах локального Set недостаточно.

Используется Redis:

const exists = await redis.get(message.id);

if (exists) {
    return;
}

await redis.set(message.id, 1, {
    EX: 300
});

Exactly Once Delivery

Большинство брокеров не гарантируют истинную модель exactly once.

На практике обычно используются:

  • at most once;
  • at least once.

Именно поэтому защита от дублей почти всегда реализуется на прикладном уровне.


At Most Once

Сообщение:

  • либо доставляется один раз;
  • либо теряется.

Повторов практически нет.

Но возможна потеря данных.


At Least Once

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

Но возможны:

  • повторы;
  • redelivery;
  • повторные ACK cycle.

Именно этот режим чаще всего используется совместно со STOMP.


Типичная архитектура защиты от дублей

Надёжная система обычно включает:

  1. UUID сообщений.
  2. ACK после успешной обработки.
  3. Идемпотентную бизнес-логику.
  4. Ограниченный кеш сообщений.
  5. Redis-дедупликацию.
  6. Контроль reconnect.
  7. Централизованный менеджер подписок.
  8. Cleanup подписок при destroy/unmount.
  9. Singleton-клиент.
  10. Логирование subscription lifecycle.

Логирование подписок

Очень полезно отслеживать создание подписок.

function createSubscription(destination, callback) {

    console.log(
        'SUBSCRIBE:',
        destination,
        Date.now()
    );

    return client.subscribe(destination, callback);
}

Отладка активных подписок

console.log(client);

Многие реализации STOMP.js содержат внутренние структуры:

  • active subscriptions;
  • reconnect state;
  • websocket state;
  • heartbeat timers.

Это помогает обнаруживать накопление подписок.


Диагностика через DevTools

Вкладка Network позволяет увидеть:

  • количество WebSocket-соединений;
  • количество reconnect;
  • повторные SUBSCRIBE frames;
  • повторные MESSAGE frames.

Если в WebSocket visible multiple SUBSCRIBE frames для одного destination — проблема практически всегда связана с повторной подпиской.


Использование debug

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);
    }
}

Такой подход:

  • предотвращает накопление подписок;
  • исключает повторные subscribe;
  • контролирует reconnect;
  • обеспечивает централизованную дедупликацию подписок;
  • упрощает диагностику проблем.