Отложенная доставка сообщений — механизм, при котором сообщение публикуется в брокер, но становится доступным подписчикам только через заданный промежуток времени или в определённый момент. В экосистеме STOMP.js этот функционал реализуется не самой библиотекой, а возможностями брокера сообщений, с которым работает клиент.
STOMP.js выступает транспортным клиентом и передаёт специальные заголовки, управляющие временем доставки.
Наиболее распространённые сценарии:
STOMP.js не содержит встроенного таймера доставки сообщений. Архитектура выглядит следующим образом:
STOMP.js → STOMP Broker → Delayed Queue/Exchange → Consumer
Клиент:
SEND.Брокер:
RabbitMQ поддерживает delayed delivery несколькими способами:
Наиболее удобный вариант — плагин
rabbitmq_delayed_message_exchange.
Пример отправки сообщения через STOMP.js:
import { Client } fr om '@stomp/stompjs';
const client = new Client({
brokerURL: 'ws://localhost:15674/ws'
});
client.onConn ect = () => {
client.publish({
destination: '/exchange/delayed-exchange/tasks',
body: JSON.stringify({
task: 'generate-report'
}),
headers: {
'x-delay': '10000'
}
});
};
client.activate();
Заголовок:
x-delay: 10000
означает задержку в 10 секунд.
После получения сообщения RabbitMQ:
Преимущества подхода:
Альтернативный подход — использование TTL.
Схема:
Producer
↓
Delay Queue
↓ (TTL expired)
Dead Letter Exchange
↓
Target Queue
↓
Consumer
Пример отправки:
client.publish({
destination: '/queue/delay-queue',
body: JSON.stringify({
id: 55
}),
headers: {
expiration: '30000'
}
});
Здесь:
expiration: 30000
означает TTL сообщения 30 секунд.
После истечения времени сообщение будет перенаправлено через DLX.
Особенности:
Недостаток:
Если первое сообщение имеет TTL 5 минут, а второе — 5 секунд, второе сообщение не будет обработано раньше первого.
Особенности:
Недостаток:
ActiveMQ поддерживает scheduled messages.
Для этого используются специальные заголовки:
AMQ_SCHEDULED_DELAY
AMQ_SCHEDULED_PERIOD
AMQ_SCHEDULED_REPEAT
Пример:
client.publish({
destination: '/queue/orders',
body: JSON.stringify({
orderId: 1001
}),
headers: {
'AMQ_SCHEDULED_DELAY': '15000'
}
});
Сообщение будет доставлено через 15 секунд.
ActiveMQ позволяет создавать повторяющиеся доставки.
Пример:
client.publish({
destination: '/queue/reminders',
body: 'ping',
headers: {
'AMQ_SCHEDULED_DELAY': '5000',
'AMQ_SCHEDULED_PERIOD': '10000',
'AMQ_SCHEDULED_REPEAT': '5'
}
});
Поведение:
Некоторые брокеры поддерживают cron-выражения.
Пример для ActiveMQ:
client.publish({
destination: '/queue/cron',
body: 'night-task',
headers: {
'AMQ_SCHEDULED_CRON': '0 0 2 * * ?'
}
});
Задача будет выполняться ежедневно в 02:00.
Одна из ключевых задач отложенной доставки — повторная обработка ошибок.
Типичный сценарий:
1 попытка → ошибка
↓
retry через 5 секунд
↓
ошибка
↓
retry через 30 секунд
↓
ошибка
↓
retry через 5 минут
Пример:
function retryPublish(client, payload, retry) {
const delay = Math.min(
1000 * Math.pow(2, retry),
60000
);
client.publish({
destination: '/exchange/retry/tasks',
body: JSON.stringify({
...payload,
retry
}),
headers: {
'x-delay': String(delay)
}
});
}
Используется exponential backoff.
Пример задержек:
| Попытка | Задержка |
|---|---|
| 1 | 1 сек |
| 2 | 2 сек |
| 3 | 4 сек |
| 4 | 8 сек |
| 5 | 16 сек |
Отложенная доставка тесно связана с DLQ.
Структура:
Main Queue
↓
Error
↓
Retry Queue
↓
Delay
↓
Main Queue
После превышения лимита попыток сообщение отправляется в DLQ.
При retry обычно добавляются заголовки:
headers: {
'x-retry-count': '3',
'x-original-queue': 'orders',
'x-last-error': 'timeout'
}
Это помогает:
При использовании delayed delivery важно правильно обрабатывать подтверждения.
Ошибочный сценарий:
1. Consumer получил сообщение
2. Произошла ошибка
3. ACK уже отправлен
4. Сообщение потеряно
Правильный сценарий:
1. Consumer получил сообщение
2. Произошла ошибка
3. Сообщение публикуется в retry queue
4. ACK отправляется только после retry publish
client.subscribe('/queue/tasks', async message => {
try {
const data = JSON.parse(message.body);
await processTask(data);
message.ack();
} catch (error) {
const retry = Number(
message.headers['x-retry-count'] || 0
);
client.publish({
destination: '/exchange/retry/tasks',
body: message.body,
headers: {
'x-delay': String((retry + 1) * 5000),
'x-retry-count': String(retry + 1)
}
});
message.ack();
}
}, {
ack: 'client'
});
Отложенная доставка увеличивает вероятность дублирования сообщений.
Причины:
Поэтому обработчики должны быть идемпотентными.
Типичный подход:
client.publish({
destination: '/queue/payments',
body: JSON.stringify({
paymentId: 'pay-100'
}),
headers: {
'message-id': crypto.randomUUID()
}
});
Получатель сохраняет обработанные идентификаторы.
Отложенная доставка позволяет реализовать планировщик задач.
Пример задач:
Проблема массовых delayed messages:
100 000 сообщений
↓
одновременное истечение TTL
↓
нагрузочный пик
Такое состояние называется burst delivery.
Вместо одинакового delay используется jitter.
Пример:
function withJitter(baseDelay) {
const random = Math.random() * 5000;
return baseDelay + random;
}
Отправка:
client.publish({
destination: '/exchange/tasks',
body: payload,
headers: {
'x-delay': String(withJitter(10000))
}
});
Важно отслеживать:
retry_count / total_messages
delivery_time - publish_time
messages_ready
dead_letter_messages
Отложенная доставка чувствительна к времени системы.
Проблемы:
Если брокеры работают в кластере, требуется единая синхронизация времени.
Некоторые брокеры плохо работают с большими delay.
Проблемные сценарии:
delay = 30 дней
delay = 365 дней
Причины:
Для очень долгих задач обычно используют:
Большинство брокеров не поддерживает удаление delayed message напрямую.
Распространённые решения:
{
id: 'task-1',
cancelled: true
}
Consumer игнорирует задачу.
Перед выполнением consumer проверяет состояние задачи:
const task = await repository.find(id);
if (task.cancelled) {
return;
}
Ошибочная конфигурация может создать retry storm.
Пример:
Ошибка
↓
retry
↓
ошибка
↓
retry
↓
бесконечный цикл
Обязательно ограничивается максимальное число попыток:
if (retry >= 5) {
client.publish({
destination: '/queue/dlq',
body: message.body
});
message.ack();
return;
}
STOMP поддерживает транзакции.
Это важно для атомарности:
publish retry
+
ack original
const tx = client.begin();
try {
client.publish({
destination: '/exchange/retry',
body: message.body,
headers: {
transaction: tx.id,
'x-delay': '5000'
}
});
message.ack({
transaction: tx.id
});
tx.commit();
} catch (error) {
tx.abort();
}
Большое количество delayed messages может приводить к:
Особенно это заметно при миллионах таймеров.
Плохой вариант:
delay = 180 дней
Лучше:
короткие retry + внешний scheduler
Рекомендуется:
3–10 попыток
Фиксированная задержка создаёт повторные пики нагрузки.
Это снижает вероятность synchronized retry storms.
Сообщения не должны теряться после превышения лимита retry.
Это критично для distributed systems.
Delayed queues способны незаметно накапливать миллионы сообщений.
STOMP.js Producer
↓
RabbitMQ Delayed Exchange
↓
Main Queue
↓
Consumer
↓
Ошибка?
↙ ↘
Нет Да
↓ ↓
ACK Retry Exchange
↓
Delayed Delivery
↓
Main Queue
↓
retry > lim it
↓
DLQ