Кеширование сообщений

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

В системах реального времени кеширование особенно важно при:

  • повторном отображении одинаковых данных;
  • временной потере соединения;
  • повторной отправке сообщений;
  • синхронизации состояния интерфейса;
  • восстановлении после реконнекта;
  • обработке burst-нагрузки;
  • оптимизации подписок.

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


Причины появления избыточных сообщений

Без кеширования приложение может:

  • многократно запрашивать одинаковые данные;
  • заново обрабатывать уже полученные события;
  • повторно отправлять идентичные команды;
  • создавать нагрузку на память и CPU;
  • увеличивать задержки интерфейса.

Типичный пример — чат:

client.subscribe('/topic/chat', (message) => {
    renderMessage(JSON.parse(message.body));
});

Если история сообщений уже была загружена ранее, повторная обработка каждого сообщения после реконнекта приводит к:

  • дублированию DOM-элементов;
  • лишнему рендерингу;
  • увеличению использования памяти;
  • ухудшению производительности.

Виды кеширования

Кеш входящих сообщений

Хранение ранее полученных сообщений.

Пример:

const messageCache = new Map();

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

    messageCache.set(data.id, data);

    processOrder(data);
});

Кеш исходящих сообщений

Используется при временном отсутствии соединения.

const pendingMessages = [];

function sendOrder(order) {
    const payload = JSON.stringify(order);

    if (client.connected) {
        client.publish({
            destination: '/app/orders',
            body: payload
        });
    } else {
        pendingMessages.push(payload);
    }
}

После восстановления соединения сообщения отправляются повторно.


Кеш подтверждений ACK

Используется для контроля обработки сообщений.

const ackCache = new Set();

client.subscribe('/queue/tasks', (message) => {
    if (ackCache.has(message.headers['message-id'])) {
        return;
    }

    ackCache.add(message.headers['message-id']);

    processTask(message);

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

Кеш состояний интерфейса

Позволяет избежать лишней перерисовки.

const stateCache = {};

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

    if (JSON.stringify(stateCache) === JSON.stringify(state)) {
        return;
    }

    Object.assign(stateCache, state);

    renderDashboard(state);
});

Кеширование при реконнекте

Во время реконнекта приложение может потерять часть сообщений. Один из способов минимизации потерь — локальный буфер.

const offlineQueue = [];

client.onDisconn ect = () => {
    console.log('Disconnected');
};

client.onConn ect = () => {
    while (offlineQueue.length > 0) {
        const message = offlineQueue.shift();

        client.publish(message);
    }
};

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

Для долговременного кеширования можно использовать LocalStorage.

Сохранение сообщений

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

    const cached = JSON.parse(
        localStorage.getItem('news-cache') || '[]'
    );

    cached.push(data);

    localStorage.setItem(
        'news-cache',
        JSON.stringify(cached)
    );
});

Восстановление кеша

const cachedMessages = JSON.parse(
    localStorage.getItem('news-cache') || '[]'
);

cachedMessages.forEach(renderNews);

IndexedDB для большого объёма сообщений

LocalStorage плохо подходит для больших объёмов данных:

  • синхронные операции;
  • ограничение размера;
  • блокировка main thread.

Для крупных realtime-систем лучше использовать IndexedDB.

Пример инициализации

const request = indexedDB.open('stomp-cache', 1);

request.onupgradenee ded = (event) => {
    const db = event.target.result;

    db.createObjectStore('messages', {
        keyPath: 'id'
    });
};

Сохранение сообщения

function saveMessage(db, message) {
    const tx = db.transaction('messages', 'readwrite');

    tx.objectStore('messages').put(message);
}

Ограничение размера кеша

Без ограничения кеш быстро увеличивается.

FIFO-очистка

const MAX_CACHE_SIZE = 1000;

const cache = [];

function addToCache(message) {
    cache.push(message);

    if (cache.length > MAX_CACHE_SIZE) {
        cache.shift();
    }
}

LRU-кеш

LRU удаляет редко используемые элементы.

class LRUCache {
    constructor(limit = 100) {
        this.limit = limit;
        this.cache = new Map();
    }

    get(key) {
        if (!this.cache.has(key)) {
            return null;
        }

        const value = this.cache.get(key);

        this.cache.delete(key);
        this.cache.set(key, value);

        return value;
    }

    set(key, value) {
        if (this.cache.has(key)) {
            this.cache.delete(key);
        }

        this.cache.set(key, value);

        if (this.cache.size > this.limit) {
            const oldest = this.cache.keys().next().value;

            this.cache.delete(oldest);
        }
    }
}

Дедупликация сообщений

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

Проверка по message-id

const processedMessages = new Set();

client.subscribe('/topic/events', (message) => {
    const id = message.headers['message-id'];

    if (processedMessages.has(id)) {
        return;
    }

    processedMessages.add(id);

    handleEvent(message);
});

TTL кеша

Иногда сообщения должны автоматически удаляться спустя определённое время.

class TTLCache {
    constructor(ttl) {
        this.ttl = ttl;
        this.cache = new Map();
    }

    set(key, value) {
        this.cache.set(key, {
            value,
            timestamp: Date.now()
        });
    }

    get(key) {
        const item = this.cache.get(key);

        if (!item) {
            return null;
        }

        if (Date.now() - item.timestamp > this.ttl) {
            this.cache.delete(key);

            return null;
        }

        return item.value;
    }
}

Кеширование батчей сообщений

При высокой нагрузке выгодно сохранять сообщения группами.

const batch = [];

client.subscribe('/topic/logs', (message) => {
    batch.push(JSON.parse(message.body));

    if (batch.length >= 100) {
        flushBatch();
    }
});

function flushBatch() {
    console.log('Saving batch', batch.length);

    batch.length = 0;
}

Комбинирование кеша и throttle

Иногда поток сообщений слишком интенсивен.

let latestMessage = null;

client.subscribe('/topic/metrics', (message) => {
    latestMessage = JSON.parse(message.body);
});

setInterval(() => {
    if (!latestMessage) {
        return;
    }

    renderMetrics(latestMessage);

    latestMessage = null;
}, 1000);

Кеширование сессионных данных

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

function getCacheKey(userId, topic) {
    return `${userId}:${topic}`;
}

Версионирование кеша

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

const CACHE_VERSION = 'v2';

function saveCache(data) {
    localStorage.setItem(
        `cache:${CACHE_VERSION}`,
        JSON.stringify(data)
    );
}

Частичное обновление кеша

Полная замена объекта не всегда эффективна.

Неэффективный вариант

cache[user.id] = user;

Частичное обновление

Object.assign(cache[user.id], {
    online: user.online
});

Иммутабельный кеш

Подход часто используется в React/Vue.

const nextState = {
    ...prevState,
    messages: [
        ...prevState.messages,
        newMessage
    ]
};

Кеширование подписок

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

const activeSubscriptions = new Map();

function subscribeOnce(destination, handler) {
    if (activeSubscriptions.has(destination)) {
        return;
    }

    const subscription = client.subscribe(
        destination,
        handler
    );

    activeSubscriptions.set(destination, subscription);
}

Сериализация кеша

Иногда требуется экспорт состояния приложения.

function exportCache(cache) {
    return JSON.stringify(cache);
}

function importCache(data) {
    return JSON.parse(data);
}

Кеширование бинарных данных

STOMP.js может работать с бинарными payload.

client.subscribe('/topic/binary', (message) => {
    const bytes = message.binaryBody;

    binaryCache.push(bytes);
});

Проблемы роста памяти

Распространённые причины утечек:

  • бесконечный массив сообщений;
  • отсутствие очистки Set/Map;
  • хранение больших бинарных объектов;
  • дублирование данных;
  • отсутствие TTL.

Плохой пример

const cache = [];

client.subscribe('/topic/data', (message) => {
    cache.push(message);
});

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


Очистка кеша

Полная очистка

cache.clear();

Частичная очистка

for (const [key, value] of cache.entries()) {
    if (value.expired) {
        cache.delete(key);
    }
}

Кеширование и производительность UI

Большое количество сообщений способно перегрузить интерфейс.

Буферизация рендера

const renderQueue = [];

client.subscribe('/topic/feed', (message) => {
    renderQueue.push(
        JSON.parse(message.body)
    );
});

setInterval(() => {
    if (renderQueue.length === 0) {
        return;
    }

    renderFeed(renderQueue.splice(0));
}, 500);

Архитектура централизованного кеша

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

class MessageCacheService {
    constructor() {
        this.messages = new Map();
    }

    save(message) {
        this.messages.set(message.id, message);
    }

    get(id) {
        return this.messages.get(id);
    }

    remove(id) {
        this.messages.delete(id);
    }

    clear() {
        this.messages.clear();
    }
}

Интеграция кеша с Redux

function reducer(state = initialState, action) {
    switch (action.type) {
        case 'MESSAGE_RECEIVED':
            return {
                ...state,
                messages: {
                    ...state.messages,
                    [action.payload.id]: action.payload
                }
            };

        default:
            return state;
    }
}

Стратегии обновления кеша

Cache Aside

Сначала проверяется кеш.

if (cache.has(id)) {
    return cache.get(id);
}

Write Through

Запись одновременно в кеш и постоянное хранилище.

cache.set(id, data);

database.save(data);

Write Back

Сначала запись в кеш, затем отложенная синхронизация.

cache.set(id, data);

syncQueue.push(data);

Мониторинг кеша

Важно отслеживать:

  • размер кеша;
  • количество попаданий;
  • количество промахов;
  • время жизни объектов;
  • скорость очистки;
  • рост памяти.

Пример метрик

const metrics = {
    hits: 0,
    misses: 0
};

function getCached(id) {
    if (cache.has(id)) {
        metrics.hits++;

        return cache.get(id);
    }

    metrics.misses++;

    return null;
}

Практическая схема кеширования realtime-приложения

Типичная архитектура:

  1. STOMP.js получает сообщения.
  2. Данные проходят дедупликацию.
  3. Сообщения попадают в memory cache.
  4. Важные данные сохраняются в IndexedDB.
  5. UI получает только изменённые состояния.
  6. Старые записи удаляются через TTL.
  7. При реконнекте выполняется повторная синхронизация.

Подобная схема позволяет:

  • снизить сетевую нагрузку;
  • уменьшить количество рендеров;
  • сократить использование CPU;
  • повысить устойчивость приложения;
  • минимизировать потери данных;
  • ускорить восстановление соединения;
  • стабилизировать работу интерфейса под высокой нагрузкой.