Кеширование сообщений в STOMP.js используется для временного хранения данных, которые уже были получены или подготовлены к отправке. Основная цель — уменьшение количества сетевых операций, снижение нагрузки на брокер сообщений и повышение отзывчивости приложения.
В системах реального времени кеширование особенно важно при:
STOMP.js сам по себе не предоставляет встроенную систему кеширования сообщений уровня приложения, однако библиотека позволяет эффективно реализовать такие механизмы поверх WebSocket и STOMP-протокола.
Без кеширования приложение может:
Типичный пример — чат:
client.subscribe('/topic/chat', (message) => {
renderMessage(JSON.parse(message.body));
});
Если история сообщений уже была загружена ранее, повторная обработка каждого сообщения после реконнекта приводит к:
Хранение ранее полученных сообщений.
Пример:
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);
}
}
После восстановления соединения сообщения отправляются повторно.
Используется для контроля обработки сообщений.
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.
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);
LocalStorage плохо подходит для больших объёмов данных:
Для крупных 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);
}
Без ограничения кеш быстро увеличивается.
const MAX_CACHE_SIZE = 1000;
const cache = [];
function addToCache(message) {
cache.push(message);
if (cache.length > MAX_CACHE_SIZE) {
cache.shift();
}
}
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);
}
}
}
Брокеры могут повторно доставлять сообщения.
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);
});
Иногда сообщения должны автоматически удаляться спустя определённое время.
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;
}
Иногда поток сообщений слишком интенсивен.
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);
});
Распространённые причины утечек:
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);
}
}
Большое количество сообщений способно перегрузить интерфейс.
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();
}
}
function reducer(state = initialState, action) {
switch (action.type) {
case 'MESSAGE_RECEIVED':
return {
...state,
messages: {
...state.messages,
[action.payload.id]: action.payload
}
};
default:
return state;
}
}
Сначала проверяется кеш.
if (cache.has(id)) {
return cache.get(id);
}
Запись одновременно в кеш и постоянное хранилище.
cache.set(id, data);
database.save(data);
Сначала запись в кеш, затем отложенная синхронизация.
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;
}
Типичная архитектура:
Подобная схема позволяет: