Пошук уроків, статей та іншого контенту
Покаже, як у RabbitMQ налаштовувати обміни, черги, прив’язки, маршрутизацію та підтвердження повідомлень.
RabbitMQ передає повідомлення через кілька основних компонентів:
Producer — створює та публікує повідомлення.
Exchange — приймає повідомлення й визначає, до яких черг їх передати.
Binding — правило зв’язку між exchange та queue.
Queue — зберігає повідомлення до моменту їх обробки.
Consumer — отримує повідомлення з черги та обробляє їх.
Producer зазвичай не надсилає повідомлення безпосередньо в чергу. Він публікує повідомлення в exchange, а exchange маршрутизує його до однієї або кількох черг.
Producer → Exchange → Binding → Queue → ConsumerExchange не є сховищем повідомлень. Він лише виконує маршрутизацію. Якщо повідомлення не потрапило до жодної черги, воно може бути відкинуте, якщо під час публікації не використано спеціальний режим повернення повідомлень.
Exchange має тип, який визначає правила маршрутизації.
direct порівнює routing key повідомлення з binding key черги. Повідомлення потрапляє до черги, якщо ключі збігаються точно.
Наприклад:
routing key: order.created
binding key: order.createdТаке повідомлення буде доставлено до черги.
Застосування:
події конкретного типу;
команди для певного обробника;
маршрутизація за точним ключем.
topic підтримує шаблони в binding key:
* — рівно одне слово;
# — нуль або більше слів.
Слова в routing key розділяються крапками.
Наприклад, для routing key:
order.eu.createdможуть підходити такі binding key:
order.*.created
order.#
#.createdАле order.* не підійде, оскільки шаблон очікує лише два слова, а в ключі їх три.
Topic exchange зручний для ієрархічних подій:
user.created
user.deleted
order.eu.created
order.us.createdfanout надсилає повідомлення в усі прив’язані черги. Routing key ігнорується.
Такий exchange підходить для широкомовних подій:
оновлення кешу;
сповіщення кількох незалежних сервісів;
розсилка події всім підписникам.
Кожен сервіс має використовувати власну чергу. Якщо кілька consumers читають одну чергу, вони розподіляють повідомлення між собою, а не отримують кожен копію.
Черга зберігає повідомлення, які ще не були успішно оброблені.
Під час створення черги важливо визначити її властивості.
Durable-черга переживає перезапуск RabbitMQ. Це не означає, що будь-яке повідомлення автоматично буде збережене на диску. Для цього повідомлення також потрібно публікувати як persistent.
await channel.assertQueue("orders", {
durable: true
});Для надійної доставки зазвичай використовують разом:
durable exchange;
durable queue;
persistent повідомлення;
підтвердження обробки consumer-ом.
Exclusive-черга належить одному з’єднанню й автоматично видаляється після його закриття. Вона корисна для тимчасових підписок, але не для основних робочих черг.
Така черга видаляється, коли з неї забирають останнього consumer-а. Її також зазвичай використовують для тимчасових підписок.
Для постійних обробників найчастіше потрібна durable-черга зі стабільним іменем.
Binding з’єднує exchange із чергою та містить правило маршрутизації.
Для direct:
await channel.bindQueue(
"orders",
"events",
"order.created"
);Для topic:
await channel.bindQueue(
"europe-orders",
"events",
"order.eu.*"
);Для fanout ключ не має значення:
await channel.bindQueue(
"cache-invalidator",
"broadcast",
""
);Exchange і queue потрібно оголосити до створення binding. Операції assertExchange, assertQueue та bindQueue є ідемпотентними: повторний виклик із такими самими параметрами не створює дублікати.
Producer публікує повідомлення в exchange, передаючи routing key.
У Node.js для роботи з RabbitMQ часто використовують пакет amqplib.
npm install amqplibПриклад публікації:
channel.publish(
"events",
"order.eu.created",
Buffer.from(JSON.stringify({
orderId: "ord-1001",
country: "EU"
})),
{
persistent: true,
contentType: "application/json"
}
);persistent: true просить RabbitMQ зберігати повідомлення на диску, якщо це можливо. Але це не замінює підтвердження від RabbitMQ про прийняття повідомлення. Для критичних producer-ів використовують publisher confirms через confirm channel.
Consumer може отримувати повідомлення в одному з двох режимів:
auto-ack — RabbitMQ вважає повідомлення обробленим одразу після доставки;
manual ack — consumer явно підтверджує успішну обробку.
Для важливих повідомлень краще використовувати manual acknowledgements:
channel.consume("orders", async (message) => {
try {
await processMessage(message);
channel.ack(message);
} catch (error) {
channel.nack(message, false, true);
}
}, {
noAck: false
});Методи:
ack(message) — повідомлення успішно оброблено;
nack(message, false, true) — обробка не вдалася, повернути повідомлення в чергу;
nack(message, false, false) — відхилити повідомлення без повторної постановки в чергу;
reject(message, true) — старіший варіант відхилення одного повідомлення з повторною постановкою.
Якщо consumer завершився після отримання повідомлення, але до виклику ack, RabbitMQ може доставити це повідомлення знову.
Тому обробник має бути готовим до повторної доставки. Наприклад, створення замовлення зазвичай повинно бути ідемпотентним: повторне отримання тієї самої події не має створити два замовлення.
nack із повторною постановкою:
channel.nack(message, false, true);може спричинити нескінченний цикл, якщо повідомлення завжди некоректне. Для таких випадків у production-системах зазвичай налаштовують окрему чергу для невдалих повідомлень, але її конфігурація виходить за межі цього прикладу.
За замовчуванням RabbitMQ може доставити consumer-у кілька повідомлень наперед. prefetch обмежує кількість непідтверджених повідомлень:
channel.prefetch(1);Значення 1 означає, що consumer отримає наступне повідомлення лише після підтвердження попереднього.
Це корисно, коли:
обробка повідомлення тривала;
важливо рівномірно розподіляти роботу між workers;
один consumer не повинен накопичувати багато непідтверджених повідомлень.
Занадто мале значення може зменшити пропускну здатність, тому його підбирають відповідно до характеру обробки.
У прикладі нижче:
створюється topic exchange;
черга europe-orders отримує європейські замовлення;
черга all-orders отримує всі замовлення;
повідомлення обробляються вручну;
після успішної обробки викликається ack;
у разі помилки повідомлення повертається в чергу.
Перед запуском має бути доступний RabbitMQ, наприклад на localhost:5672.
const amqp = require("amqplib");
const RABBITMQ_URL = "amqp://localhost:5672";
const EXCHANGE = "orders.events";
async function main() {
const connection = await amqp.connect(RABBITMQ_URL);
const channel = await connection.createChannel();
await channel.assertExchange(EXCHANGE, "topic", {
durable: true
});
await channel.assertQueue("europe-orders", {
durable: true
});
await channel.assertQueue("all-orders", {
durable: true
});
await channel.bindQueue(
"europe-orders",
EXCHANGE,
"order.eu.*"
);
await channel.bindQueue(
"all-orders",
EXCHANGE,
"order.#"
);
channel.prefetch(1);
await channel.consume(
"europe-orders",
async (message) => {
if (!message) {
return;
}
try {
const event = JSON.parse(message.content.toString());
console.log(
"[europe-orders] Обробка:",
event
);
await processOrder(event);
channel.ack(message);
} catch (error) {
console.error(
"[europe-orders] Помилка:",
error.message
);
// Повертаємо повідомлення в чергу для повторної спроби.
channel.nack(message, false, true);
}
},
{
noAck: false
}
);
await channel.consume(
"all-orders",
async (message) => {
if (!message) {
return;
}
try {
const event = JSON.parse(message.content.toString());
console.log(
"[all-orders] Отримано:",
event
);
await processOrder(event);
channel.ack(message);
} catch (error) {
console.error(
"[all-orders] Помилка:",
error.message
);
// У цьому прикладі помилкове повідомлення буде доставлене повторно.
channel.nack(message, false, true);
}
},
{
noAck: false
}
);
publishOrder(channel, "order.eu.created", {
orderId: "ord-1001",
region: "eu"
});
publishOrder(channel, "order.us.created", {
orderId: "ord-1002",
region: "us"
});
console.log("Consumers запущені. Натисніть Ctrl+C для завершення.");
}
function publishOrder(channel, routingKey, order) {
const accepted = channel.publish(
EXCHANGE,
routingKey,
Buffer.from(JSON.stringify(order)),
{
persistent: true,
contentType: "application/json"
}
);
if (!accepted) {
console.warn(
"Внутрішній буфер каналу переповнений"
);
}
}
async function processOrder(order) {
if (!order.orderId) {
throw new Error("У повідомленні немає orderId");
}
await new Promise((resolve) => {
setTimeout(resolve, 300);
});
console.log(
`Замовлення ${order.orderId} успішно оброблено`
);
}
main().catch((error) => {
console.error("Не вдалося запустити застосунок:", error);
process.exitCode = 1;
});Для повідомлення з ключем order.eu.created RabbitMQ знайде обидва binding:
order.eu.*
order.#Тому повідомлення потрапить і до europe-orders, і до all-orders.
Для повідомлення з ключем order.us.created підійде лише:
order.#Тому воно потрапить тільки до all-orders.
Це важлива властивість RabbitMQ: одне повідомлення може бути скопійоване до кількох черг, якщо воно відповідає binding-ам цих черг.
Consumer acknowledgements відповідають на питання:
Чи успішно consumer обробив повідомлення?
Окремо існують publisher confirms, які відповідають на питання:
Чи RabbitMQ прийняв опубліковане повідомлення?
Для publisher confirms створюють confirm channel:
const confirmChannel = await connection.createConfirmChannel();
await confirmChannel.assertExchange("events", "direct", {
durable: true
});
confirmChannel.publish(
"events",
"order.created",
Buffer.from(JSON.stringify({
orderId: "ord-1003"
})),
{
persistent: true,
contentType: "application/json"
}
);
await confirmChannel.waitForConfirms();
console.log("RabbitMQ підтвердив публікацію");Це не є підтвердженням бізнес-обробки. RabbitMQ підтверджує прийняття повідомлення, а ack від consumer-а — завершення його обробки. Це різні етапи доставки.
noAck: trueЯкщо повідомлення потрібно обробити надійно, не використовуйте автоматичне підтвердження:
{
noAck: true
}У такому разі RabbitMQ вважатиме повідомлення обробленим ще до виконання бізнес-логіки. При падінні застосунку повідомлення вже не буде доставлене повторно.
ack до завершення обробкиНеправильно:
channel.consume("orders", async (message) => {
channel.ack(message);
await saveOrder(message);
});Якщо saveOrder завершиться помилкою, RabbitMQ вже видалив повідомлення з черги.
Підтвердження потрібно викликати після успішного завершення всіх необхідних операцій.
Producer передає routing key під час публікації. Binding key належить зв’язку між exchange та queue.
Для direct ключі мають збігатися точно. Для topic binding key може містити * і #.
Якщо два сервіси читають одну чергу, RabbitMQ розподіляє повідомлення між ними. Другий сервіс не отримає копію кожного повідомлення.
Якщо кожен сервіс має отримати всі події, створіть для кожного сервісу окрему чергу та прив’яжіть її до exchange.
Exchange потрібно оголосити перед публікацією. Якщо producer публікує в неіснуючий exchange, RabbitMQ закриє канал із помилкою.
Durable-черга зберігає свою конфігурацію після перезапуску RabbitMQ. Для збереження повідомлень також потрібно публікувати їх із persistent: true і враховувати підтвердження від RabbitMQ.
Безумовний виклик:
channel.nack(message, false, true);для невиправної помилки може нескінченно повертати те саме повідомлення в чергу.
Потрібно розрізняти:
тимчасову помилку, після якої можна повторити обробку;
постійну помилку даних, яку повторна спроба не виправить.
Producer публікує повідомлення в exchange, а не безпосередньо в queue.
Exchange маршрутизує повідомлення через binding.
direct використовує точний збіг routing key.
topic підтримує шаблони * і #.
fanout надсилає повідомлення в усі прив’язані черги.
Durable-черга переживає перезапуск RabbitMQ, а persistent повідомлення призначене для збереження на диску.
Для надійної обробки використовуйте manual acknowledgements.
ack викликається після успішної обробки.
nack може підтвердити помилку та повернути повідомлення в чергу або відхилити його.
prefetch обмежує кількість непідтверджених повідомлень у consumer-а.
Consumer acknowledgements і publisher confirms підтверджують різні етапи доставки.