Ограничение частоты запросов (Rate Limiting) — механизм контроля количества сообщений, отправляемых клиентом через STOMP-соединение за определённый промежуток времени. В контексте STOMP.js этот механизм особенно важен при работе с:
Без ограничения частоты клиент способен:
Даже один ошибочный цикл setInterval() способен
генерировать тысячи сообщений в секунду.
Если сообщения отправляются быстрее, чем брокер успевает их обрабатывать, очередь начинает стремительно расти.
Пример проблемы:
setInterval(() => {
client.publish({
destination: '/topic/data',
body: JSON.stringify(generateData())
});
}, 1);
Такой код создаёт до 1000 сообщений в секунду.
STOMP.js работает поверх WebSocket. При слишком высокой скорости передачи возникают:
Ошибки фронтенда могут выглядеть как атака:
while (true) {
client.publish({
destination: '/topic/logs',
body: 'spam'
});
}
Некоторые брокеры автоматически разрывают соединение.
STOMP.js временно хранит данные в памяти JavaScript-движка. Если отправка быстрее сети, формируется внутренний буфер.
Последствия:
Фиксированное окно времени.
Например:
Схема:
0s ------- 1s ------- 2s
| 100 req | 100 req |
Недостаток — всплески на границе окна.
Скользящее окно анализирует последние N миллисекунд.
Пример:
Последние 1000 ms → максимум 100 сообщений
Механизм более точный и плавный.
Одна из лучших стратегий для STOMP.
Принцип:
Имитирует постоянную скорость утечки.
Даже если клиент создаёт всплеск, сообщения выходят равномерно.
Подходит для:
import { Client } fr om '@stomp/stompjs';
const client = new Client({
brokerURL: 'ws://localhost:15674/ws'
});
let lastSend = 0;
const LIMIT_MS = 1000;
function safePublish(message) {
const now = Date.now();
if (now - lastSend < LIMIT_MS) {
console.warn('Лимит превышен');
return;
}
lastSend = now;
client.publish({
destination: '/topic/messages',
body: JSON.stringify(message)
});
}
Подобный вариант:
let counter = 0;
setInterval(() => {
counter = 0;
}, 1000);
function publishLimited(body) {
if (counter >= 10) {
console.warn('Слишком много сообщений');
return;
}
counter++;
client.publish({
destination: '/topic/chat',
body
});
}
Преимущества:
Недостатки:
Пусть:
class TokenBucket {
constructor(maxTokens, refillRate) {
this.maxTokens = maxTokens;
this.tokens = maxTokens;
this.refillRate = refillRate;
setInterval(() => {
this.tokens = Math.min(
this.maxTokens,
this.tokens + this.refillRate
);
}, 1000);
}
consume(count = 1) {
if (this.tokens < count) {
return false;
}
this.tokens -= count;
return true;
}
}
const bucket = new TokenBucket(20, 5);
function publish(body) {
if (!bucket.consume()) {
console.warn('Rate lim it exceeded');
return;
}
client.publish({
destination: '/topic/events',
body
});
}
Если клиент простаивал:
Это удобно для:
Иногда нельзя терять сообщения.
Вместо:
return;
лучше помещать сообщение в очередь.
class PublishQueue {
constructor(ratePerSecond) {
this.queue = [];
this.interval = 1000 / ratePerSecond;
this.start();
}
enqueue(message) {
this.queue.push(message);
}
start() {
setInterval(() => {
if (this.queue.length === 0) {
return;
}
const message = this.queue.shift();
client.publish({
destination: message.destination,
body: message.body
});
}, this.interval);
}
}
const queue = new PublishQueue(5);
queue.enqueue({
destination: '/topic/chat',
body: 'hello'
});
Очередь обеспечивает:
Debounce откладывает выполнение до завершения серии вызовов.
Пример:
function debounce(fn, delay) {
let timeout;
return (...args) => {
clearTimeout(timeout);
timeout = setTimeout(() => {
fn(...args);
}, delay);
};
}
const sendTyping = debounce(() => {
client.publish({
destination: '/topic/typing',
body: 'typing'
});
}, 500);
Подходит для:
Throttle ограничивает частоту выполнения.
function throttle(fn, limit) {
let waiting = false;
return (...args) => {
if (waiting) {
return;
}
fn(...args);
waiting = true;
setTimeout(() => {
waiting = false;
}, limit);
};
}
const sendMouse = throttle((position) => {
client.publish({
destination: '/topic/mouse',
body: JSON.stringify(position)
});
}, 100);
Особенно эффективен для:
Rate limiting касается не только количества сообщений.
Важно ограничивать:
function publishSafe(body) {
const bytes = new Blob([body]).size;
if (bytes > 1024 * 10) {
throw new Error('Payload too large');
}
client.publish({
destination: '/topic/data',
body
});
}
На практике используется несколько уровней защиты одновременно.
Пример:
При падении брокера тысячи клиентов начинают переподключение одновременно.
Это создаёт:
client.reconnectDelay = 100;
Слишком агрессивное переподключение.
Правильнее увеличивать задержку постепенно.
let reconnectDelay = 1000;
client.onWebSocketCl ose = () => {
setTimeout(() => {
client.activate();
reconnectDelay *= 2;
reconnectDelay = Math.min(reconnectDelay, 30000);
}, reconnectDelay);
};
Некоторые приложения динамически создают подписки.
Ошибка:
setInterval(() => {
client.subscribe('/topic/random', () => {});
}, 10);
Это вызывает:
const subscriptions = new Map();
const MAX_SUBSCRIPTIONS = 50;
function safeSubscribe(destination, callback) {
if (subscriptions.size >= MAX_SUBSCRIPTIONS) {
throw new Error('Subscription limit exceeded');
}
const sub = client.subscribe(destination, callback);
subscriptions.set(destination, sub);
return sub;
}
Клиентский код можно:
Настоящий контроль выполняется на сервере.
В RabbitMQ используются:
ActiveMQ поддерживает:
Spring позволяет:
Подозрительные признаки:
const history = [];
function trackMessage() {
const now = Date.now();
history.push(now);
while (history[0] < now - 1000) {
history.shift();
}
if (history.length > 100) {
console.warn('Flood detected');
}
}
RxJS хорошо подходит для потокового ограничения.
import { Subject } from 'rxjs';
import { throttleTime } from 'rxjs/operators';
const stream = new Subject();
stream
.pipe(throttleTime(100))
.subscribe(message => {
client.publish({
destination: '/topic/data',
body: message
});
});
import { bufferTime } from 'rxjs/operators';
stream
.pipe(bufferTime(1000))
.subscribe(messages => {
client.publish({
destination: '/topic/batch',
body: JSON.stringify(messages)
});
});
Грамотно настроенное ограничение:
Крупные системы обычно используют:
Подходящие лимиты:
Рекомендуется:
Подходят:
Важно: