Пошук уроків, статей та іншого контенту
Розкриє ролі продюсерів і консьюмерів, формат повідомлень, підтвердження обробки та керування навантаженням.
У системах обміну повідомленнями одна частина програми створює повідомлення, а інша їх отримує та обробляє.
Продюсер — компонент, який створює і надсилає повідомлення.
Консьюмер — компонент, який отримує та обробляє повідомлення.
Брокер повідомлень — посередник, який приймає повідомлення від продюсерів і передає їх консьюмерам.
Наприклад, після оформлення замовлення інтернет-магазин може не виконувати всі дії одразу:
сервіс замовлень створює замовлення;
продюсер надсилає повідомлення order.created;
брокер зберігає це повідомлення;
консьюмер надсилає лист клієнту;
інший консьюмер резервує товар на складі.
Завдяки цьому сервіси не повинні напряму викликати один одного.
Спрощена схема виглядає так:
Продюсер → Брокер повідомлень → КонсьюмерПродюсер не обов’язково знає:
де саме працює консьюмер;
скільки консьюмерів існує;
коли консьюмер завершить обробку;
чи тимчасово недоступний консьюмер.
Брокер відповідає за доставку повідомлень і часто зберігає їх доти, доки консьюмер не підтвердить успішну обробку.
Продюсер може надіслати таке повідомлення:
{
"type": "order.created",
"messageId": "msg-8f31",
"createdAt": "2026-09-02T10:15:00.000Z",
"payload": {
"orderId": "order-123",
"customerId": "customer-42",
"total": 1499.99
}
}Консьюмер отримує повідомлення і, наприклад, надсилає клієнту підтвердження замовлення.
Повідомлення зазвичай має дві частини:
метадані — інформація для інфраструктури й обробки;
корисне навантаження (payload) — дані предметної області.
До метаданих можуть належати:
messageId — унікальний ідентифікатор повідомлення;
type — тип події або команди;
createdAt — час створення;
correlationId — ідентифікатор пов’язаного запиту чи операції;
attempt — номер спроби обробки.
У різних системах назви полів можуть відрізнятися, але призначення залишається подібним.
payload містить дані, потрібні конкретному консьюмеру:
{
"orderId": "order-123",
"customerId": "customer-42",
"total": 1499.99
}Не варто без потреби додавати в повідомлення великі об’єкти або конфіденційні дані. Повідомлення має містити рівно стільки інформації, скільки потрібно для його обробки.
Повідомлення часто поділяють на події та команди.
Подія описує те, що вже сталося:
order.created
payment.completed
user.registeredПодія не вимагає від конкретного консьюмера виконати дію. Вона повідомляє про факт.
Наприклад:
{
"type": "payment.completed",
"payload": {
"paymentId": "payment-77",
"orderId": "order-123"
}
}Команда просить виконати певну дію:
send.welcome.email
reserve.product
generate.invoiceНаприклад:
{
"type": "send.welcome.email",
"payload": {
"userId": "user-15",
"email": "user@example.com"
}
}Для початку важливо розуміти головне: продюсер формує повідомлення, а консьюмер використовує його для виконання роботи.
Після отримання повідомлення консьюмер має повідомити брокеру, що обробка завершилася успішно. Таке повідомлення називають підтвердженням або acknowledgement (ack).
Типовий порядок дій:
консьюмер отримує повідомлення;
виконує необхідну роботу;
перевіряє результат;
надсилає ack.
Брокер → повідомлення → Консьюмер
Брокер ← ack ← КонсьюмерЯкщо консьюмер аварійно завершився до надсилання ack, брокер може повторно доставити повідомлення.
Це захищає повідомлення від втрати, але створює іншу вимогу: обробка має бути готовою до повторного запуску.
Якщо обробка успішна:
отримати повідомлення → обробити → ackЯкщо сталася тимчасова помилка:
отримати повідомлення → помилка → повторити пізнішеЯкщо повідомлення неправильне або обробити його неможливо:
отримати повідомлення → помилка → перемістити в окрему чергуТаку окрему чергу часто називають dead-letter queue. Вона допомагає не блокувати обробку інших повідомлень і зберегти проблемне повідомлення для аналізу.
Повідомлення може бути доставлене більше одного разу. Наприклад:
консьюмер оновив дані в базі;
перед надсиланням ack процес завершився;
брокер не отримав підтвердження;
те саме повідомлення було доставлено повторно.
Якщо консьюмер повторить операцію без перевірки, можуть виникнути проблеми:
клієнт отримає два однакові листи;
гроші спишуться двічі;
товар зарезервується кілька разів.
Тому обробку часто роблять ідемпотентною. Це означає, що повторна обробка того самого повідомлення не змінює правильний результат.
Один із простих підходів:
отримати messageId;
перевірити, чи оброблялося це повідомлення раніше;
якщо так — не виконувати роботу повторно;
якщо ні — виконати роботу й зберегти факт обробки.
messageId = msg-8f31
якщо msg-8f31 уже оброблено:
надіслати ack
інакше:
виконати операцію
зберегти msg-8f31 як оброблене
надіслати ackПродюсер може створювати повідомлення швидше, ніж консьюмер встигає їх обробляти.
Наприклад:
продюсер надсилає 1000 повідомлень за секунду;
один консьюмер обробляє лише 100 повідомлень за секунду.
У такому разі черга зростатиме. Якщо не контролювати навантаження, система може використати забагато пам’яті або збільшити час очікування.
Консьюмер може отримувати лише обмежену кількість повідомлень, наприклад 10. Нові повідомлення він бере тільки після завершення попередніх.
Це називають обмеженням кількості повідомлень «у польоті».
максимум одночасних повідомлень = 10Якщо повідомлень стабільно більше, ніж може обробити один консьюмер, можна запустити кілька його екземплярів.
Черга → Консьюмер 1
→ Консьюмер 2
→ Консьюмер 3Повідомлення розподіляються між ними, тому загальна пропускна здатність зростає.
Продюсер може тимчасово зменшити швидкість створення повідомлень, якщо система перевантажена. Такий підхід називають backpressure — поширенням сигналу про перевантаження назад до продюсера.
Нижче наведено спрощену чергу повідомлень. Вона показує:
роботу продюсера;
роботу консьюмера;
обмеження кількості одночасних повідомлень;
підтвердження після успішної обробки.
class MessageQueue {
constructor() {
this.messages = [];
this.waitingConsumers = [];
}
publish(message) {
const consumer = this.waitingConsumers.shift();
if (consumer) {
consumer(message);
} else {
this.messages.push(message);
}
}
consume(onMessage) {
const message = this.messages.shift();
if (message) {
onMessage(message);
} else {
this.waitingConsumers.push(onMessage);
}
}
size() {
return this.messages.length;
}
}
const queue = new MessageQueue();
let activeMessages = 0;
const maxConcurrentMessages = 2;
function publishOrder(orderId) {
const message = {
messageId: `message-${orderId}`,
type: "order.created",
createdAt: new Date().toISOString(),
payload: {
orderId,
total: 1499.99
}
};
queue.publish(message);
console.log(`Продюсер: надіслано ${message.messageId}`);
}
function startConsumer() {
if (activeMessages >= maxConcurrentMessages) {
return;
}
activeMessages += 1;
queue.consume(async (message) => {
try {
console.log(`Консьюмер: отримано ${message.messageId}`);
await new Promise((resolve) => setTimeout(resolve, 500));
console.log(`Консьюмер: замовлення ${message.payload.orderId} оброблено`);
console.log(`Консьюмер: ack для ${message.messageId}`);
} finally {
activeMessages -= 1;
startConsumer();
}
});
}
for (let i = 1; i <= 5; i += 1) {
publishOrder(`order-${i}`);
}
const timer = setInterval(() => {
startConsumer();
if (queue.size() === 0 && activeMessages === 0) {
clearInterval(timer);
console.log("Усі повідомлення оброблено");
}
}, 100);У цьому прикладі:
publishOrder виконує роль продюсера;
MessageQueue виконує спрощену роль брокера;
startConsumer запускає обробку повідомлень;
maxConcurrentMessages обмежує кількість одночасних обробок;
ack моделюється повідомленням у консолі після успішної роботи.
Реальні брокери мають додаткові можливості: збереження повідомлень, повторну доставку, кілька черг і груп консьюмерів. Але основна модель залишається такою самою.
Кілька консьюмерів можуть працювати разом як група.
У такому випадку кожне повідомлення обробляє один консьюмер із групи:
Черга:
повідомлення 1 → Консьюмер A
повідомлення 2 → Консьюмер B
повідомлення 3 → Консьюмер AЦе дає змогу розподілити навантаження.
Водночас різні групи можуть отримувати власні копії повідомлень. Наприклад:
Подія order.created
├─ група email-сервісу
└─ група складського сервісуОбидві групи можуть незалежно обробити одну й ту саму подію для різних цілей.
Повний цикл можна описати так:
продюсер створює повідомлення;
повідомлення отримує унікальний ідентифікатор;
продюсер надсилає його брокеру;
брокер зберігає повідомлення в черзі;
консьюмер отримує повідомлення;
консьюмер виконує бізнес-операцію;
консьюмер надсилає ack;
брокер видаляє повідомлення або позначає його обробленим.
Якщо під час обробки виникла помилка, можливий інший сценарій:
консьюмер отримує повідомлення;
обробка завершується помилкою;
ack не надсилається;
брокер повторює доставку або переносить повідомлення в окрему чергу.
Сам факт отримання повідомлення ще не означає, що операція завершилася успішно. Підтвердження потрібно надсилати після завершення необхідної роботи.
ack до виконання операціїЯкщо консьюмер спочатку надсилає ack, а потім падає під час обробки, брокер може більше не доставити повідомлення. У результаті дані буде втрачено.
Консьюмер має бути готовим отримати те саме повідомлення кілька разів. Для цього використовують messageId і перевірку вже оброблених повідомлень.
Якщо консьюмер отримує більше повідомлень, ніж може обробити, зростає використання пам’яті та час очікування. Потрібно встановлювати розумне обмеження паралельних операцій.
Великі повідомлення повільніше передавати й зберігати. Краще передавати необхідні дані або ідентифікатор ресурсу, який консьюмер може отримати окремо.
Тимчасову помилку, наприклад короткочасну недоступність сервісу, можна повторити. Некоректне повідомлення повторювати безкінечно не має сенсу — його варто перемістити в окрему чергу для аналізу.
Продюсер створює та надсилає повідомлення.
Консьюмер отримує повідомлення й виконує роботу.
Брокер зберігає та доставляє повідомлення.
Повідомлення зазвичай містить метадані й payload.
ack підтверджує успішне завершення обробки.
Відсутність ack може призвести до повторної доставки.
Обробка повідомлень має бути готовою до повторного запуску.
Навантаження контролюють обмеженням паралельної роботи, масштабуванням консьюмерів і backpressure.
Унікальний messageId допомагає відстежувати повідомлення та уникати повторного виконання операцій.