Отложенная доставка

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

STOMP.js выступает транспортным клиентом и передаёт специальные заголовки, управляющие временем доставки.

Наиболее распространённые сценарии:

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

Архитектура отложенной доставки

STOMP.js не содержит встроенного таймера доставки сообщений. Архитектура выглядит следующим образом:

STOMP.js → STOMP Broker → Delayed Queue/Exchange → Consumer

Клиент:

  1. Формирует сообщение.
  2. Указывает специальные заголовки.
  3. Отправляет фрейм SEND.

Брокер:

  1. Принимает сообщение.
  2. Помещает его во временное хранилище.
  3. Удерживает до наступления времени доставки.
  4. Перенаправляет подписчикам.

Отложенная доставка в RabbitMQ

RabbitMQ поддерживает delayed delivery несколькими способами:

  • TTL + Dead Letter Exchange;
  • delayed message exchange plugin;
  • quorum queues с retry;
  • delayed queues через x-message-ttl.

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

Наиболее удобный вариант — плагин 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 секунд.


Как работает x-delay

После получения сообщения RabbitMQ:

  1. Сохраняет сообщение.
  2. Вычисляет момент публикации.
  3. Удерживает сообщение внутри exchange.
  4. Отправляет в очередь после истечения задержки.

Преимущества подхода:

  • отсутствуют промежуточные retry-очереди;
  • простая конфигурация;
  • минимальный объём инфраструктуры;
  • высокая читаемость архитектуры.

Отложенная доставка через TTL

Альтернативный подход — использование 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 и x-delay

TTL

Особенности:

  • сообщение сначала попадает в очередь;
  • используется dead-letter routing;
  • может блокироваться head-of-line эффектом;
  • сложнее конфигурировать.

Недостаток:

Если первое сообщение имеет TTL 5 минут, а второе — 5 секунд, второе сообщение не будет обработано раньше первого.


x-delay

Особенности:

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

Недостаток:

  • требуется RabbitMQ plugin.

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

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

Поведение:

  • первая отправка через 5 секунд;
  • затем повтор каждые 10 секунд;
  • всего 5 повторов.

Планирование по cron

Некоторые брокеры поддерживают cron-выражения.

Пример для ActiveMQ:

client.publish({
    destination: '/queue/cron',
    body: 'night-task',
    headers: {
        'AMQ_SCHEDULED_CRON': '0 0 2 * * ?'
    }
});

Задача будет выполняться ежедневно в 02:00.


Retry-механизмы через delayed delivery

Одна из ключевых задач отложенной доставки — повторная обработка ошибок.

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

1 попытка → ошибка
↓
retry через 5 секунд
↓
ошибка
↓
retry через 30 секунд
↓
ошибка
↓
retry через 5 минут

Реализация retry в STOMP.js

Пример:

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 сек

Dead Letter Queues

Отложенная доставка тесно связана с DLQ.

Структура:

Main Queue
    ↓
Error
    ↓
Retry Queue
    ↓
Delay
    ↓
Main Queue

После превышения лимита попыток сообщение отправляется в DLQ.


Метаданные повторных попыток

При retry обычно добавляются заголовки:

headers: {
    'x-retry-count': '3',
    'x-original-queue': 'orders',
    'x-last-error': 'timeout'
}

Это помогает:

  • анализировать ошибки;
  • строить monitoring;
  • реализовывать circuit breaker;
  • отслеживать проблемные сообщения.

Отложенная доставка и ACK

При использовании delayed delivery важно правильно обрабатывать подтверждения.

Ошибочный сценарий:

1. Consumer получил сообщение
2. Произошла ошибка
3. ACK уже отправлен
4. Сообщение потеряно

Правильный сценарий:

1. Consumer получил сообщение
2. Произошла ошибка
3. Сообщение публикуется в retry queue
4. ACK отправляется только после retry publish

Пример безопасного retry

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

Идемпотентность сообщений

Отложенная доставка увеличивает вероятность дублирования сообщений.

Причины:

  • reconnect;
  • network split;
  • broker failover;
  • повторные публикации;
  • race conditions;
  • повторная отправка retry.

Поэтому обработчики должны быть идемпотентными.


Использование message-id

Типичный подход:

client.publish({
    destination: '/queue/payments',
    body: JSON.stringify({
        paymentId: 'pay-100'
    }),
    headers: {
        'message-id': crypto.randomUUID()
    }
});

Получатель сохраняет обработанные идентификаторы.


Scheduled tasks через STOMP.js

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

Пример задач:

  • отправка email;
  • генерация отчётов;
  • очистка кэша;
  • резервное копирование;
  • напоминания;
  • webhooks;
  • аналитические вычисления.

Массовое планирование сообщений

Проблема массовых 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))
    }
});

Мониторинг delayed delivery

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

  • размер delayed queues;
  • скорость доставки;
  • количество retry;
  • DLQ growth;
  • latency;
  • backlog;
  • delivery lag.

Полезные метрики

Retry rate

retry_count / total_messages

Delivery latency

delivery_time - publish_time

Queue depth

messages_ready

DLQ size

dead_letter_messages

Временная синхронизация

Отложенная доставка чувствительна к времени системы.

Проблемы:

  • рассинхронизация серверов;
  • некорректный NTP;
  • drift виртуальных машин;
  • проблемы контейнерных сред.

Если брокеры работают в кластере, требуется единая синхронизация времени.


Долгие задержки

Некоторые брокеры плохо работают с большими delay.

Проблемные сценарии:

delay = 30 дней
delay = 365 дней

Причины:

  • рост памяти;
  • огромные индексы таймеров;
  • медленный recovery;
  • увеличение snapshot-файлов.

Для очень долгих задач обычно используют:

  • внешние scheduler-сервисы;
  • базы данных;
  • cron-системы;
  • workflow engines.

Отмена отложенных сообщений

Большинство брокеров не поддерживает удаление delayed message напрямую.

Распространённые решения:

Флаг отмены

{
    id: 'task-1',
    cancelled: true
}

Consumer игнорирует задачу.


Хранилище состояния

Перед выполнением consumer проверяет состояние задачи:

const task = await repository.find(id);

if (task.cancelled) {
    return;
}

Защита от бесконечных retry

Ошибочная конфигурация может создать retry storm.

Пример:

Ошибка
↓
retry
↓
ошибка
↓
retry
↓
бесконечный цикл

Обязательно ограничивается максимальное число попыток:

if (retry >= 5) {

    client.publish({
        destination: '/queue/dlq',
        body: message.body
    });

    message.ack();

    return;
}

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

STOMP поддерживает транзакции.

Это важно для атомарности:

publish retry
+
ack original

Пример transaction

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 может приводить к:

  • росту RAM;
  • увеличению IO;
  • замедлению recovery;
  • деградации broker throughput;
  • росту latency;
  • увеличению времени failover.

Особенно это заметно при миллионах таймеров.


Практические рекомендации

Не использовать огромные delay

Плохой вариант:

delay = 180 дней

Лучше:

короткие retry + внешний scheduler

Ограничивать retry

Рекомендуется:

3–10 попыток

Использовать exponential backoff

Фиксированная задержка создаёт повторные пики нагрузки.


Добавлять jitter

Это снижает вероятность synchronized retry storms.


Использовать DLQ

Сообщения не должны теряться после превышения лимита retry.


Делать обработчики идемпотентными

Это критично для distributed systems.


Отслеживать backlog

Delayed queues способны незаметно накапливать миллионы сообщений.


Пример полноценной архитектуры

STOMP.js Producer
        ↓
RabbitMQ Delayed Exchange
        ↓
Main Queue
        ↓
Consumer
        ↓
Ошибка?
   ↙         ↘
Нет           Да
 ↓             ↓
ACK        Retry Exchange
                ↓
          Delayed Delivery
                ↓
           Main Queue
                ↓
        retry > lim it
                ↓
               DLQ