Пошук уроків, статей та іншого контенту
Додасте retry, backoff, dead-letter черги та контроль помилок у розподілених сервісах.
У розподіленій системі помилка не завжди означає, що операцію потрібно негайно завершити. Наприклад:
зовнішній сервіс тимчасово недоступний;
з'єднання з базою даних було розірвано;
повідомлення обробляється довше, ніж дозволяє тайм-аут;
сервіс повернув тимчасову помилку 5xx.
У таких випадках повторна спроба може завершити операцію успішно. Проте без обмежень retry здатен створити нескінченний цикл і перевантажити систему.
Типова стратегія містить:
обмежену кількість повторних спроб;
затримку між спробами;
збільшення затримки після кожної невдалої спроби;
окрему dead-letter чергу для повідомлень, які не вдалося обробити;
розділення тимчасових і постійних помилок.
За фіксованої затримки кожна повторна спроба виконується через однаковий проміжок часу:
спроба 1 → 5 секунд → спроба 2 → 5 секунд → спроба 3Цей варіант простий, але багато повідомлень можуть повторно надходити одночасно й створювати нове навантаження.
Exponential backoff збільшує затримку після кожної невдалої спроби:
затримка = базоваЗатримка × 2 ^ номерСпробиНаприклад, для базової затримки 5 секунд:
спроба 1: 5 секунд
спроба 2: 10 секунд
спроба 3: 20 секундЩоб уникнути одночасного повторного надходження великої кількості повідомлень, до затримки часто додають невеликий випадковий компонент — jitter.
function getBackoffDelay(
attempt: number,
baseDelayMs = 5_000,
maxDelayMs = 60_000,
): number {
const exponentialDelay = baseDelayMs * 2 ** attempt;
const jitter = Math.floor(Math.random() * 1_000);
return Math.min(exponentialDelay + jitter, maxDelayMs);
}У RabbitMQ затримку можна реалізувати через retry-черги з TTL. Після завершення TTL RabbitMQ перемістить повідомлення назад у головну чергу через dead-letter exchange.
Dead-letter queue, або DLQ, призначена для повідомлень, які не можна або не вдалося обробити.
Повідомлення варто направляти до DLQ, якщо:
вичерпано максимальну кількість спроб;
помилка є постійною;
повідомлення має неправильний формат;
бізнес-правило забороняє повторну обробку.
DLQ не повинна бути просто місцем, де повідомлення «зникають». Зазвичай її використовують для:
аналізу причин помилки;
ручного виправлення даних;
повторного запуску після виправлення проблеми;
моніторингу та сповіщень.
Для прикладу використаємо такі черги:
orders
├── orders.retry.1 ── 5 секунд ──┐
├── orders.retry.2 ── 10 секунд ─┤──> orders
└── orders.dlqАлгоритм:
Сервіс отримує повідомлення з orders.
Якщо обробка успішна, повідомлення підтверджується через ack.
Якщо помилка тимчасова:
збільшується лічильник спроб;
повідомлення публікується в одну з retry-черг;
оригінальне повідомлення підтверджується.
Після завершення TTL RabbitMQ повертає повідомлення в orders.
Якщо кількість спроб перевищила ліміт, повідомлення публікується в orders.dlq.
Після публікації повідомлення в retry-чергу потрібно підтвердити оригінальне повідомлення лише після успішної публікації. Інакше повідомлення можна втратити.
Встановимо залежності:
npm install @nestjs/common @nestjs/core amqplib reflect-metadata rxjs
npm install --save-dev @types/amqplib typescript ts-nodeНижче наведено приклад NestJS-сервісу, який:
створює RabbitMQ-черги;
використовує ручне підтвердження повідомлень;
виконує повторні спроби;
застосовує exponential backoff через TTL-черги;
переміщує невдалі повідомлення до DLQ.
import {
Injectable,
Logger,
Module,
OnModuleDestroy,
OnModuleInit,
} from '@nestjs/common';
import { NestFactory } from '@nestjs/core';
import {
Channel,
ConfirmChannel,
Connection,
ConsumeMessage,
Options,
} from 'amqplib';
import * as amqp from 'amqplib';
const MAIN_QUEUE = 'orders';
const RETRY_QUEUES = ['orders.retry.1', 'orders.retry.2'];
const DLQ = 'orders.dlq';
const MAX_RETRIES = RETRY_QUEUES.length;
type OrderMessage = {
orderId: string;
status?: string;
};
@Injectable()
class OrdersConsumer implements OnModuleInit, OnModuleDestroy {
private readonly logger = new Logger(OrdersConsumer.name);
private connection!: Connection;
private channel!: ConfirmChannel;
async onModuleInit(): Promise<void> {
const url = process.env.RABBITMQ_URL ?? 'amqp://localhost:5672';
this.connection = await amqp.connect(url);
this.channel = await this.connection.createConfirmChannel();
await this.setupQueues();
await this.channel.prefetch(10);
await this.channel.consume(
MAIN_QUEUE,
(message) => {
if (message) {
void this.handleMessage(message);
}
},
{ noAck: false },
);
this.logger.log(`Очікування повідомлень у черзі "${MAIN_QUEUE}"`);
}
private async setupQueues(): Promise<void> {
await this.channel.assertQueue(MAIN_QUEUE, {
durable: true,
});
await this.channel.assertQueue(RETRY_QUEUES[0], {
durable: true,
arguments: {
'x-message-ttl': 5_000,
'x-dead-letter-exchange': '',
'x-dead-letter-routing-key': MAIN_QUEUE,
},
});
await this.channel.assertQueue(RETRY_QUEUES[1], {
durable: true,
arguments: {
'x-message-ttl': 10_000,
'x-dead-letter-exchange': '',
'x-dead-letter-routing-key': MAIN_QUEUE,
},
});
await this.channel.assertQueue(DLQ, {
durable: true,
});
}
private async handleMessage(message: ConsumeMessage): Promise<void> {
const headers = message.properties.headers ?? {};
const retryCount = Number(headers['x-retry-count'] ?? 0);
try {
const payload = this.parseMessage(message);
await this.processOrder(payload);
this.channel.ack(message);
this.logger.log(`Замовлення ${payload.orderId} успішно оброблено`);
} catch (error) {
const normalizedError = this.normalizeError(error);
this.logger.error(
`Помилка обробки повідомлення. Спроба: ${retryCount}`,
normalizedError.stack,
);
try {
if (normalizedError.permanent || retryCount >= MAX_RETRIES) {
await this.publishToDeadLetterQueue(message, normalizedError);
} else {
await this.publishToRetryQueue(message, retryCount, normalizedError);
}
// Підтверджуємо оригінал лише після успішної публікації копії.
this.channel.ack(message);
} catch (publishError) {
this.logger.error(
'Не вдалося перенаправити повідомлення',
publishError instanceof Error ? publishError.stack : String(publishError),
);
// Якщо нове повідомлення не опубліковане, дозволяємо повторну доставку.
this.channel.nack(message, false, true);
}
}
}
private parseMessage(message: ConsumeMessage): OrderMessage {
try {
const value = JSON.parse(message.content.toString()) as OrderMessage;
if (!value.orderId || typeof value.orderId !== 'string') {
throw new Error('Поле orderId є обов’язковим');
}
return value;
} catch (error) {
const parseError = new Error(
`Некоректний формат повідомлення: ${
error instanceof Error ? error.message : String(error)
}`,
);
Object.assign(parseError, { permanent: true });
throw parseError;
}
}
private async processOrder(order: OrderMessage): Promise<void> {
// Імітація тимчасової помилки зовнішнього сервісу.
if (order.status === 'temporary-failure') {
const error = new Error('Зовнішній сервіс тимчасово недоступний');
Object.assign(error, {
code: 'EXTERNAL_SERVICE_UNAVAILABLE',
permanent: false,
});
throw error;
}
// Імітація постійної бізнес-помилки.
if (order.status === 'invalid') {
const error = new Error('Замовлення має недійсний статус');
Object.assign(error, {
code: 'INVALID_ORDER_STATUS',
permanent: true,
});
throw error;
}
// Тут могла б бути транзакція в базі даних або виклик іншого сервісу.
await Promise.resolve();
}
private async publishToRetryQueue(
message: ConsumeMessage,
retryCount: number,
error: NormalizedError,
): Promise<void> {
const nextRetryCount = retryCount + 1;
const retryQueue = RETRY_QUEUES[nextRetryCount - 1];
const headers = {
...(message.properties.headers ?? {}),
'x-retry-count': nextRetryCount,
'x-last-error': error.message.slice(0, 500),
};
this.channel.sendToQueue(retryQueue, message.content, {
persistent: true,
contentType: message.properties.contentType,
correlationId: message.properties.correlationId,
headers,
});
await this.channel.waitForConfirms();
this.logger.warn(
`Повідомлення перенаправлено в ${retryQueue}. ` +
`Наступна спроба: ${nextRetryCount}`,
);
}
private async publishToDeadLetterQueue(
message: ConsumeMessage,
error: NormalizedError,
): Promise<void> {
const headers = {
...(message.properties.headers ?? {}),
'x-final-error': error.message.slice(0, 500),
'x-error-code': error.code ?? 'UNKNOWN_ERROR',
};
this.channel.sendToQueue(DLQ, message.content, {
persistent: true,
contentType: message.properties.contentType,
correlationId: message.properties.correlationId,
headers,
});
await this.channel.waitForConfirms();
this.logger.error(`Повідомлення перенаправлено в ${DLQ}`);
}
private normalizeError(error: unknown): NormalizedError {
if (error instanceof Error) {
const typedError = error as Error & {
code?: string;
permanent?: boolean;
};
return {
message: typedError.message,
stack: typedError.stack,
code: typedError.code,
permanent: typedError.permanent === true,
};
}
return {
message: String(error),
permanent: false,
};
}
async onModuleDestroy(): Promise<void> {
await this.channel?.close();
await this.connection?.close();
}
}
type NormalizedError = {
message: string;
stack?: string;
code?: string;
permanent: boolean;
};
@Module({
providers: [OrdersConsumer],
})
class AppModule {}
async function bootstrap(): Promise<void> {
const app = await NestFactory.createApplicationContext(AppModule);
const shutdown = async (): Promise<void> => {
await app.close();
process.exit(0);
};
process.once('SIGINT', () => void shutdown());
process.once('SIGTERM', () => void shutdown());
}
void bootstrap();Для локального запуску RabbitMQ можна використати контейнер:
docker run --rm \
--name lesson-rabbitmq \
-p 5672:5672 \
rabbitmq:3Після запуску RabbitMQ NestJS-сервіс підключиться до нього за адресою amqp://localhost:5672.
Не кожна помилка повинна спричиняти retry.
Повторна спроба зазвичай доречна для таких ситуацій:
тайм-аут мережевого запиту;
тимчасова недоступність сервісу;
помилка підключення до бази даних;
відповідь 429 Too Many Requests;
відповідь 502, 503 або 504.
Повторення не допоможе, якщо:
JSON має неправильний формат;
відсутнє обов'язкове поле;
ідентифікатор не існує;
порушено бізнес-правило;
запит не має необхідних прав;
операція дублює вже виконану дію, а сервіс не підтримує її повторення.
Корисно створювати власні класи помилок або додавати до помилки ознаку permanent. У production-коді краще не визначати тип помилки лише за текстом повідомлення.
Наприклад:
class PermanentProcessingError extends Error {
readonly permanent = true;
readonly code: string;
constructor(message: string, code: string) {
super(message);
this.name = 'PermanentProcessingError';
this.code = code;
}
}
class TemporaryProcessingError extends Error {
readonly permanent = false;
readonly code: string;
constructor(message: string, code: string) {
super(message);
this.name = 'TemporaryProcessingError';
this.code = code;
}
}Тоді маршрутизація помилки може спиратися на тип або властивість, а не на порівняння рядків.
RabbitMQ підтримує два основні варіанти роботи зі споживачем:
autoAck: true — повідомлення вважається обробленим одразу після доставки;
noAck: false — споживач сам викликає ack або nack.
Для надійної обробки потрібен ручний режим.
ackack підтверджує успішну обробку повідомлення. Після цього RabbitMQ більше не доставлятиме його цьому споживачу.
nacknack повідомляє брокеру, що обробка завершилася помилкою.
channel.nack(message, false, true);Третій аргумент визначає, чи потрібно повернути повідомлення в чергу:
true — повідомлення буде доставлене повторно;
false — повідомлення буде відкинуте або передане до DLQ, якщо для черги налаштовано dead-letter routing.
У прикладі повідомлення спочатку публікується до потрібної retry-черги або DLQ, а потім підтверджується через ack. Це дає змогу контролювати маршрут повідомлення самостійно.
Повторна доставка можлива навіть тоді, коли використовується ручний ack. Наприклад:
обробник зберіг дані в базі;
перед викликом ack з'єднання з RabbitMQ перервалося;
RabbitMQ не отримав підтвердження;
повідомлення було доставлене повторно.
Тому обробник повинен бути ідемпотентним: повторне виконання тієї самої операції не повинно створювати неправильний результат.
Практичні підходи:
використовувати унікальний orderId або eventId;
зберігати оброблені ідентифікатори;
застосовувати унікальний індекс у базі даних;
перевіряти поточний стан сутності перед зміною;
виконувати зміни в транзакції.
Retry без ідемпотентності може створити дублікати платежів, замовлень або повідомлень.
Максимальна кількість спроб повинна бути явно визначена. Вона залежить від операції:
для швидкого HTTP-виклику може бути достатньо 2–3 спроб;
для повільної інтеграції — більше;
для незворотної операції повторення може бути взагалі заборонене.
Не варто повторювати повідомлення нескінченно. Це може:
зайняти всі ресурси споживача;
збільшити навантаження на зовнішній сервіс;
приховати реальну проблему;
затримати обробку нових повідомлень.
Retry-лічильник потрібно зберігати в заголовках повідомлення або в іншому надійному сховищі. Лічильник у пам'яті процесу буде втрачено після перезапуску сервісу.
channel.ack(message);
await processMessage(message);Якщо processMessage завершиться помилкою, RabbitMQ вже вважатиме повідомлення успішно обробленим.
Правильний порядок:
await processMessage(message);
channel.ack(message);nack з повторною доставкоюchannel.nack(message, false, true);Якщо викликати цей код для кожної помилки, повідомлення може нескінченно повертатися в ту саму чергу.
Потрібно вести лічильник спроб і після його перевищення направляти повідомлення до DLQ.
Одна черга з TTL може створити небажану поведінку: повідомлення на початку черги з великою затримкою блокуватиме повідомлення, які вже готові до доставки.
Окремі черги на кшталт retry.1, retry.2, retry.3 спрощують реалізацію різних інтервалів.
Неправильний JSON або відсутнє обов'язкове поле не виправиться через 10 секунд. Такі повідомлення потрібно одразу направляти до DLQ.
Повторна спроба може бути результатом не лише помилки бізнес-логіки, а й втрати мережевого з'єднання після фактичного виконання операції. Обробник повинен безпечно працювати з дубльованими повідомленнями.
Якщо спочатку виконати ack, а потім публікувати повідомлення до retry-черги, збій публікації призведе до втрати повідомлення.
Публікуйте повідомлення через confirm channel і виконуйте ack лише після підтвердження брокером.
Retry потрібен для тимчасових, а не для всіх помилок.
Exponential backoff зменшує навантаження під час повторних спроб.
TTL retry-черги можуть реалізувати затримку в RabbitMQ.
Після вичерпання ліміту повідомлення потрібно направляти до dead-letter черги.
Використовуйте ручний ack і підтверджуйте повідомлення лише після успішної обробки.
Під час перенаправлення спочатку підтверджуйте публікацію, а потім — оригінальне повідомлення.
Споживачі повідомлень мають бути ідемпотентними.
Retry-ліміт, типи помилок і роботу DLQ потрібно контролювати через логи та метрики.