Пошук уроків, статей та іншого контенту
Розгляне стратегії retry, dead-letter queue, ідемпотентність і збереження порядку під час збоїв.
У розподіленій системі обробник повідомлення може завершитися помилкою навіть тоді, коли саме повідомлення коректне:
база даних тимчасово недоступна;
мережевий запит завершився тайм-аутом;
сервіс-партнер тимчасово перевантажений;
виникла помилка ліміту запитів;
процес обробника перезапустився після виконання операції, але до підтвердження повідомлення.
Якщо повідомлення просто видалити після першої помилки, подія буде втрачена. Якщо повторювати спроби без обмежень, одне проблемне повідомлення може заблокувати всю систему.
Тому надійна обробка зазвичай має такі елементи:
класифікацію помилок;
обмежену кількість повторних спроб;
затримку між спробами;
dead-letter queue для повідомлень, які не вдалося обробити;
ідемпотентний обробник;
явну політику щодо порядку повідомлень.
Не кожна помилка є підставою для повторної спроби.
Повторна спроба зазвичай доречна, якщо проблема може зникнути сама:
тайм-аут;
тимчасова недоступність сервісу;
HTTP 429 Too Many Requests;
HTTP 502, 503 або 504;
короткочасна помилка підключення до бази даних;
конфлікт транзакцій, який можна повторити.
Такі помилки називають transient errors.
Повторення не виправить проблему, якщо:
повідомлення має неправильний формат;
відсутнє обов’язкове поле;
порушено бізнес-правило;
посилання на ресурс не існує;
операція заборонена для цього користувача;
версія схеми повідомлення не підтримується.
Такі повідомлення слід швидко перемістити до DLQ або окремого потоку помилок.
Іноді система не може точно визначити результат операції. Наприклад, обробник відправив платіжному сервісу запит, але втратив відповідь через мережевий збій.
У цьому випадку повторна спроба може:
безпечно повторити операцію, якщо вона ідемпотентна;
створити дубль, якщо вона неілемпотентна.
Тому для невизначених помилок ідемпотентність є обов’язковою умовою безпечного retry.
Після кожної помилки система чекає однаковий проміжок часу:
1 с → 1 с → 1 сСтратегія проста, але має недоліки:
багато споживачів можуть повторити запити одночасно;
перевантажений сервіс отримає нову хвилю запитів;
затримка не враховує тривалість проблеми.
Затримка поступово збільшується:
1 с → 2 с → 3 с → 4 сЦе краще за фіксовану затримку, але при тривалому збої збільшення може бути недостатнім.
Типова формула:
delay = min(maxDelay, baseDelay * 2^attempt)Наприклад:
100 мс → 200 мс → 400 мс → 800 мс → 1600 мсПотрібно обмежувати затримку зверху, інакше остання спроба може відбутися через дуже довгий час.
Якщо всі споживачі використовують однакову формулу, вони можуть повторити запити одночасно. Це називається thundering herd problem.
Jitter додає випадкову складову:
delay = random(0, min(maxDelay, baseDelay * 2^attempt))Або:
delay = min(maxDelay, baseDelay * 2^attempt) + random(0, jitter)Для розподілених систем зазвичай використовують експоненційний backoff разом із jitter.
Retry повинен мати обмеження:
максимальна кількість спроб;
максимальний загальний час;
максимальна затримка між спробами.
Обмеження за кількістю спроб недостатнє саме по собі. Наприклад, п’ять спроб із затримкою в одну годину означають, що повідомлення може залишатися активним кілька годин.
Для кожного повідомлення корисно зберігати службові дані:
message_id
attempt
first_failed_at
last_failed_at
next_retry_at
last_errorЦі дані потрібні для:
моніторингу;
розслідування збоїв;
ручного повторного запуску;
визначення повідомлень, які зависли надовго.
У більшості брокерів споживач спочатку отримує повідомлення, а потім підтверджує його обробку.
Типовий порядок:
споживач отримує повідомлення;
виконує бізнес-операцію;
успішно завершує її;
надсилає acknowledgment;
брокер видаляє повідомлення або позначає його обробленим.
Якщо процес завершується між кроками 3 і 4, брокер може доставити повідомлення повторно. Це нормальна властивість моделі at-least-once delivery.
Важливий наслідок:
Повторна доставка може відбутися навіть після успішного виконання бізнес-операції.
Тому не можна вважати ack заміною ідемпотентності.
Якщо підтвердити повідомлення до виконання операції, можна втратити повідомлення. Якщо підтвердити після операції, можливі дублікати. Практична стратегія — підтверджувати після успішної обробки та робити саму обробку ідемпотентною.
Операція є ідемпотентною, якщо її повторне виконання має той самий підсумковий ефект, що й одне виконання.
Наприклад, встановлення стану:
status = "paid"можна повторити без зміни результату.
Натомість операція:
balance = balance + 100не є ідемпотентною без додаткового захисту. Її повторне виконання може двічі зарахувати гроші.
Повідомлення повинно мати стабільний ідентифікатор операції:
{
"message_id": "payment-8472",
"type": "PaymentCaptured",
"aggregate_id": "order-91",
"version": 4,
"amount": 2500
}Обробник може зберігати ідентифікатори вже виконаних повідомлень у таблиці або іншому надійному сховищі.
Спрощений алгоритм:
почати транзакцію;
перевірити, чи оброблявся message_id;
якщо так — завершити операцію без повторного побічного ефекту;
якщо ні — виконати зміну даних;
записати message_id як оброблений;
закомітити транзакцію;
підтвердити повідомлення брокеру.
Перевірка та запис ідентифікатора повинні бути атомарними. Інакше два паралельні обробники можуть одночасно перевірити відсутність ідентифікатора та обидва виконати операцію.
Якщо обробник викликає зовнішній сервіс, локального журналу може бути недостатньо. Наприклад:
обробник створює платіж у зовнішній системі;
запит успішно доходить;
відповідь втрачається;
обробник повторює запит.
Зовнішній сервіс також повинен підтримувати ідемпотency key або інший спосіб визначити дубль. Один і той самий ключ має відповідати одній логічній операції.
Dead-letter queue, або DLQ, — це місце для повідомлень, які система не змогла обробити після визначеної політики retry.
Повідомлення слід переміщати до DLQ, коли:
вичерпано максимальну кількість спроб;
минув максимальний час обробки;
виявлено постійну помилку;
повідомлення не відповідає схемі;
потрібне ручне рішення оператора.
DLQ не повинна бути просто «смітником». Для кожного повідомлення варто зберігати:
оригінальне тіло;
ідентифікатор повідомлення;
тип події;
ключ порядку або агрегат;
кількість спроб;
час першої та останньої помилки;
текст і тип останньої помилки;
назву початкової черги;
версію схеми;
ідентифікатор кореляції.
Наявність DLQ сама по собі не вирішує проблему. Потрібен операційний процес:
моніторити кількість нових повідомлень;
групувати помилки за причиною;
виправити код або дані;
перевірити повідомлення на тестовому середовищі;
повторно запустити повідомлення;
контролювати дублікати та порушення порядку.
Повторний запуск може виконуватися двома способами:
поверненням повідомлення до початкової черги;
публікацією нового повідомлення з посиланням на оригінальне.
Під час replay не слід безконтрольно повертати всі повідомлення одночасно. Це може створити новий сплеск навантаження або знову порушити порядок подій.
Порядок повідомлень може бути важливим лише в межах конкретного контексту. Наприклад:
події одного замовлення;
зміни одного рахунку;
команди одного пристрою;
події одного користувача.
Глобальний порядок для всіх повідомлень часто є дорогим або непотрібним обмеженням.
Потрібно розрізняти:
порядок публікації;
порядок доставки;
порядок початку обробки;
порядок завершення обробки;
порядок фіксації результату.
Навіть якщо брокер доставляє повідомлення по черзі, паралельні обробники можуть завершити їх у неправильному порядку.
Розглянемо події одного замовлення:
OrderCreated
PaymentCaptured
OrderShippedЯкщо PaymentCaptured тимчасово не обробляється, OrderShipped не можна виконати раніше. Інакше система може відправити неоплачене замовлення.
Це створює head-of-line blocking: перше проблемне повідомлення блокує всі наступні повідомлення того самого потоку.
Є кілька стратегій.
Поки повідомлення не буде успішно оброблене, наступні повідомлення того самого ключа не запускаються.
Переваги:
порядок зберігається;
бізнес-логіка простіша;
стан агрегату не переходить через пропущену подію.
Недоліки:
одна проблемна подія може надовго заблокувати ключ;
потрібні retry та моніторинг завислих потоків.
Після вичерпання retry система може продовжити обробку наступних повідомлень.
Перевага — більша пропускна здатність. Недолік — порядок порушується. Це допустимо лише тоді, коли:
події незалежні;
бізнес-логіка дозволяє пропуск;
наступні обробники можуть працювати без попередньої події;
існує механізм компенсації.
Якщо повідомлення розділені за ключами, можна заблокувати тільки потік одного замовлення, не зупиняючи інші замовлення.
Це зазвичай кращий компроміс:
order-91 → заблокований
order-92 → обробляється
order-93 → обробляєтьсяДля збереження порядку події одного ключа повинні потрапляти в одну послідовну область обробки:
partition = hash(order_id) % number_of_partitionsЦе дозволяє обробляти різні ключі паралельно, але послідовно обробляти події в межах одного ключа.
Однак партиціювання саме по собі не гарантує повної коректності:
зміна кількості партицій може змінити розподіл ключів;
один гарячий ключ може перевантажити одну партицію;
повторне розміщення повідомлення в окрему retry-чергу може порушити порядок;
паралельна обробка всередині партиції також може зламати порядок.
Потрібно заздалегідь визначити scope порядку: глобальний, на партицію чи на бізнес-ключ.
Окрема retry-черга зручна для відкладених повторних спроб, але може створити проблему:
Основна черга: A1 → A2 → A3
Retry: A1Якщо основна черга продовжить обробку, A2 і A3 завершаться раніше за A1.
Для ключів, де порядок обов’язковий, потрібно:
не відправляти наступні події ключа в обробку;
зберігати статус заблокованого ключа;
відновлювати повідомлення в правильній позиції;
або використовувати механізм брокера, який підтримує послідовну доставку та затримані повтори.
Після переміщення повідомлення до DLQ порядок також не відновлюється автоматично. Replay повинен враховувати номер версії або послідовності події.
Нижче наведено самодостатній приклад на JavaScript. Він демонструє:
незалежні потоки за key;
послідовну обробку повідомлень одного ключа;
експоненційний backoff із jitter;
обмеження кількості спроб;
ідемпотентність через набір оброблених ідентифікаторів;
переміщення повідомлення до DLQ;
блокування наступного повідомлення того самого ключа до завершення попереднього.
const messages = [
{ id: "order-1-1", key: "order-1", sequence: 1, type: "created" },
{ id: "order-1-2", key: "order-1", sequence: 2, type: "paid" },
{ id: "order-1-3", key: "order-1", sequence: 3, type: "shipped" },
{ id: "order-2-1", key: "order-2", sequence: 1, type: "created" },
];
const MAX_ATTEMPTS = 4;
const BASE_DELAY_MS = 20;
const MAX_DELAY_MS = 150;
const processedMessageIds = new Set();
const attempts = new Map();
const dlq = [];
const stateByKey = new Map();
function sleep(ms) {
return new Promise((resolve) => setTimeout(resolve, ms));
}
function getDelay(attempt) {
const exponentialDelay = Math.min(
MAX_DELAY_MS,
BASE_DELAY_MS * 2 ** (attempt - 1)
);
// Додаємо випадкову затримку, щоб споживачі не повторювали запит одночасно.
return Math.floor(Math.random() * exponentialDelay);
}
async function applyBusinessOperation(message) {
// Імітуємо тимчасовий збій для події paid під час перших двох спроб.
const currentAttempt = attempts.get(message.id) ?? 0;
if (message.type === "paid" && currentAttempt < 2) {
throw new Error("тимчасово недоступна платіжна система");
}
// Імітуємо постійний збій для події shipped без успішної оплати.
const currentState = stateByKey.get(message.key);
if (message.type === "shipped" && currentState !== "paid") {
throw new Error("замовлення не має стану paid");
}
stateByKey.set(message.key, message.type);
console.log(`Оброблено ${message.id}: ${message.type}`);
}
async function processMessage(message) {
if (processedMessageIds.has(message.id)) {
console.log(`Пропущено дубль ${message.id}`);
return true;
}
const attempt = (attempts.get(message.id) ?? 0) + 1;
attempts.set(message.id, attempt);
try {
await applyBusinessOperation(message);
// У реальній системі цей запис має бути атомарним із бізнес-зміною.
processedMessageIds.add(message.id);
return true;
} catch (error) {
console.log(
`Помилка ${message.id}, спроба ${attempt}: ${error.message}`
);
if (attempt >= MAX_ATTEMPTS) {
dlq.push({
...message,
failedAt: new Date().toISOString(),
attempts: attempt,
error: error.message,
});
console.log(`Переміщено до DLQ: ${message.id}`);
return false;
}
await sleep(getDelay(attempt));
return processMessage(message);
}
}
async function processKey(key, keyMessages) {
// Повідомлення одного ключа обробляються тільки послідовно.
const sortedMessages = [...keyMessages].sort(
(a, b) => a.sequence - b.sequence
);
for (const message of sortedMessages) {
const success = await processMessage(message);
if (!success) {
// Не обробляємо наступні події цього ключа, щоб не порушити порядок.
console.log(`Потік ${key} заблоковано після ${message.id}`);
break;
}
}
}
async function main() {
const messagesByKey = new Map();
for (const message of messages) {
if (!messagesByKey.has(message.key)) {
messagesByKey.set(message.key, []);
}
messagesByKey.get(message.key).push(message);
}
// Різні ключі можуть оброблятися паралельно.
await Promise.all(
[...messagesByKey.entries()].map(([key, keyMessages]) =>
processKey(key, keyMessages)
)
);
console.log("\nФінальний стан:");
console.log(Object.fromEntries(stateByKey));
console.log("\nDLQ:");
console.log(dlq);
}
main().catch((error) => {
console.error("Непередбачена помилка:", error);
process.exitCode = 1;
});У цьому прикладі подія paid спочатку завершується помилкою, але успішно обробляється під час повторної спроби. Подія shipped виконується після paid, тому порядок для order-1 зберігається.
Якщо повідомлення потрапляє до DLQ, наступні повідомлення цього самого ключа не обробляються. Це свідомий вибір на користь коректності порядку. Для іншої доменної моделі можна вибрати продовження обробки, але таке рішення потрібно явно зафіксувати.
Ці властивості часто плутають.
Ідемпотентність захищає від дублювання:
A → A → AПісля повторів результат залишається правильним.
Порядок захищає від неправильної послідовності:
A → B → Cзамість:
B → A → CМожлива система, яка:
зберігає порядок, але не захищає від дублікатів;
захищає від дублікатів, але обробляє події не по порядку;
має обидві властивості;
не має жодної.
Для більшості фінансових операцій і змін стану агрегатів потрібні обидві властивості.
Потрібно вимірювати не лише кількість успішно оброблених повідомлень.
Корисні метрики:
кількість повідомлень на кожній спробі;
частка повідомлень, успішних після retry;
кількість повідомлень у DLQ;
вік найстарішого повідомлення в DLQ;
час від першої помилки до успішної обробки;
кількість заблокованих ключів;
тривалість блокування ключа;
кількість повторних доставок;
розподіл помилок за типами.
Логи повинні містити стабільні поля:
message_id
key
sequence
attempt
correlation_id
error_typeБез message_id та attempt дублікати й повторні спроби складно відрізнити від нових повідомлень.
Без обмеження повідомлення з постійною помилкою може нескінченно споживати ресурси.
Як виправити: встановити максимальну кількість спроб, максимальний час і DLQ.
Це може створити синхронні хвилі повторних запитів.
Як виправити: використовувати exponential backoff із jitter.
Неправильний JSON або відсутнє поле не виправляться після десяти повторів.
Як виправити: класифікувати помилки та відразу відправляти некоректні повідомлення до DLQ.
Після тайм-ауту обробник не знає, чи завершилася операція. Повторний запуск може двічі списати кошти або створити два ресурси.
Як виправити: використовувати унікальний ідемпотency key та атомарно зберігати факт обробки.
Навіть упорядкована доставка не допоможе, якщо A і B обробляються паралельно та B завершується раніше.
Як виправити: серіалізувати обробку в межах ключа.
Переміщення A до DLQ та обробка B може зламати стан агрегату.
Як виправити: визначити, чи є порядок обов’язковим, і блокувати наступні повідомлення для залежного ключа.
Повідомлення накопичуються, але ніхто не знає, як їх безпечно повернути в систему.
Як виправити: зберігати причину помилки, версію повідомлення та мати контрольований процес повторного запуску.
Retry потрібен для тимчасових помилок, але має бути обмеженим.
Для розподілених систем типовою стратегією є exponential backoff із jitter.
Постійні помилки не слід повторювати без кінця — їх потрібно переміщати до DLQ.
DLQ повинна містити достатньо метаданих для аналізу та контрольованого replay.
Модель at-least-once означає, що дублікати можливі навіть після успішної операції.
Ідемпотентність захищає від повторної доставки, але не гарантує порядок.
Для збереження порядку повідомлення одного бізнес-ключа слід обробляти послідовно.
Помилка одного повідомлення може блокувати наступні повідомлення цього ключа.
Потрібно явно вибрати політику: блокувати ключ, пропускати повідомлення або використовувати компенсацію.
Retry, DLQ, ідемпотентність і порядок мають бути частиною єдиного контракту обробки повідомлень.