Как отправить отложенное сообщение в RabbitMQ: примеры кода

Иногда в приложение нужно добавить функцию, которая выполняет действие с задержкой или повторяется с определённым интервалом. Например, отправить пуш-уведомление через 10 минут или каждый день очищать временную папку.
Чтобы решить такую задачу, можно использовать cron-задачи — автоматический запуск скриптов на сервере в определённое время — или пакет node-schedule (библиотека для планирования задач в стиле cron на Node.js).
Но при использовании этих решений возникнет проблема с масштабированием:
- Если серверов или запущенных сервисов много, может быть непонятно, на каком именно запускать эту задачу
- Выбранный сервер может выйти из строя
- Нода может удалиться из-за освобождения ресурсов
Одно из возможных решений этой проблемы — брокер сообщений RabbitMQ. Схема реализации следующая:
1. Создаем 2 обменника: обычный и delayed
export const HELLO_EXCHANGE = Object.freeze({
name: 'hello',
type: 'direct',
options: {
durable: true,
},
queues: {},
});
export const HELLO_DELAYED_EXCHANGE = Object.freeze({
name: 'helloDelayed',
type: 'direct',
options: {
durable: true,
},
queues: {},
});
2. Создаем в них очереди с одинаковым типом связывания данных (binding), но разными именами (name)
Для HELLO_EXCHANGE:
queues: {
WORLD: {
name: 'hello.world', // subscribe to this queue
binding: 'hello.world',
options: {
durable: true,
},
},
},
Для HELLO_DELAYED_EXCHANGE:
queues: {
WORLD: {
name: 'helloDelayed.world',
binding: 'hello.world',
options: {
durable: true,
queueMode: 'lazy', // set the message to remain in the hard memory
},
}
Для очереди delayed-обменника указываем аргумент x-dead-letter-exchange с именем обычной очереди. Аргумент указывает брокеру перенести сообщение в этот обменник, если сообщение не обработается.
arguments: {
'x-dead-letter-exchange': HELLO_EXCHANGE.name, // указываем в какую очередь должно переместиться сообщение после своей смерти
}
3. Публикуем сообщение в очередь delayed-обменника с указанием срока действия — expiration
// services/base-service/src/broker/hello/publisher.ts
export const publishHelloDelayedWorld = createPublisher({
exchangeName: exchangeNameDelayed,
queue: WORLD_DELAYED,
expirationInMs: 30000, //указываем, через сколько сообщение должно умереть (30s)
});
По истечении срока expiration сообщение автоматически перенесется в очередь обычного обменника.
Теперь осталось только создать подписчика (consumer) на очередь обычного обменника:
// services/base-service/src/broker/hello/consumer.ts
export const initHelloExchange = () => Promise.all([
createConsumer(
{
queueName: HELLO_EXCHANGE.queues.WORLD.name,
prefetch: 50,
log: true,
},
controller.consumeHelloWorld,
),
]);
// services/base-service/src/broker/hello/controller.ts
export const consumeHelloWorld: IBrokerHandler = async ({ payload }) => {
const result = await world({ name: payload.name });
logger.info(result.message);
// await publishHelloDelayedWorld({ name: payload.name }); // если необходимо вновь обработать сообщение
};
Profit!
Если нужно выполнять действие периодически, в конце consumer снова отправляет сообщение в delayed-очередь.
// await publishHelloDelayedWorld({ name: payload.name });
NOTE: RabbitMQ работает по принципу FIFO (first in, first out) — команды обрабатываются в том порядке, в котором были отправлены. То есть, если сначала отправить в одну очередь сообщение с временем жизни в 1 день, а потом — ещё одно с временем жизни в 1 минуту, второе сообщение начнёт обрабатываться только после первого. Целевое действие для него выполнится через минуту после завершения обработки первого.
В конечном итоге получаем такую цепочку:
- Создаём обменники и очереди:
// services/base-service/src/broker/const/exchanges.ts
export const HELLO_EXCHANGE = Object.freeze({
name: 'hello',
type: 'direct',
options: {
durable: true,
},
queues: {
WORLD: {
name: 'hello.world', // подписываемся на эту очередь
binding: 'hello.world',
options: {
durable: true,
},
},
},
});
export const HELLO_DELAYED_EXCHANGE = Object.freeze({
name: 'helloDelayed',
type: 'direct',
options: {
durable: true,
queueMode: 'lazy', // указываем, что сообщение должно сохраниться в жесткой памяти
},
queues: {
WORLD: {
name: 'helloDelayed.world',
binding: 'hello.world',
options: {
durable: true,
queueMode: 'lazy', // указываем чтобы сообщение хранилось в жесткой памяти arguments: {
'x-dead-letter-exchange': HELLO_EXCHANGE.name, // указываем в какую очередь должно переместиться сообщение после своей смерти
},
},
},
},
});
- Добавляем publisher, который будет отправлять сообщение в задержанную очередь:
// services/base-service/src/broker/hello/publisher.ts
export const publishHelloDelayedWorld = createPublisher({
exchangeName: exchangeNameDelayed,
queue: WORLD_DELAYED,
expirationInMs: 30000, // set when the message dies (in 30s)
});
- Добавляем consumer для очереди из обычного обменника:
// services/base-service/src/broker/hello/consumer.ts
export const initHelloExchange = () => Promise.all([
createConsumer(
{
queueName: HELLO_EXCHANGE.queues.WORLD.name,
prefetch: 50,
log: true,
},
controller.consumeHelloWorld,
),
]);
// services/base-service/src/broker/hello/controller.ts
export const consumeHelloWorld: IBrokerHandler = async ({ payload }) => {
const result = await world({ name: payload.name });
logger.info(result.message);
// await publishHelloDelayedWorld({ name: payload.name }); // if you need to process the message again
};
Profit!
Есть плагин, который сделает эту работу за нас — реализация в нём проще. Нам нужно создать всего один обменник, одну очередь, один паблишер и один консьюмер.
При публикации плагин самостоятельно обработает сообщение и доставит его в нужную очередь по истечении времени.
При реализации через плагин сообщения обрабатываются по мере истечения задержки — delay. То есть, если сначала опубликовать сообщение с delay в 1 день, а потом ещё одно с delay в 1 минуту, второе сообщение будет обработано раньше первого.
// services/base-service/src/broker/const/exchanges.ts
export const HELLO_PLUGIN_DELAYED_EXCHANGE = Object.freeze({
name: 'helloPluginDelayed',
type: 'x-delayed-message', // specify the delayed queue
options: {
durable: true,
arguments: {
'x-delayed-type': 'direct', // set the recipient },
},
queues: {
WORLD_PLUGIN_DELAYED: {
name: 'helloPluginDelayed.world', // subscribe to the queue
binding: 'helloPluginDelayed.world',
options: {
durable: true,
},
},
},
});
Добавляем publisher, который будет отправлять сообщение в delayed-очередь:
export const publishHelloPluginDelayedWorld = createPublisher({
exchangeName: exchangeNamePluginDelayed,
queue: WORLD_PLUGIN_DELAYED,
delayInMs: 60000, // specify when the message should die (60s)
});
Добавляем consumer для очереди:
// services/base-service/src/broker/hello/consumer.ts
export const initHelloExchange = () => Promise.all([
createConsumer(
{
queueName: HELLO_PLUGIN_DELAYED_EXCHANGE.queues.WORLD_PLUGIN_DELAYED.name,
prefetch: 50,
log: true,
},
controller.consumeHelloWorld,
),
]);
// services/base-service/src/broker/hello/controller.ts
export const consumeHelloWorld: IBrokerHandler = async ({ payload }) => {
const result = await world({ name: payload.name });
logger.info(result.message);
};
Посмотреть реализацию можно в репозитории.
Мы уже применяли это во многих проектах. Например, в Janson Media internet TV. Это прокат фильмов, но в цифровом формате.
В этом проекте мы использовали RabbitMQ для реализации трёх основных функций сервиса: отправка SMS и писем пользователям с напоминанием об окончании срока аренды; отправка уведомлений об оплате через сокет и показ соответствующего уведомления пользователю; передача загруженных видео на дальнейшую обработку.
