Симптом механической ошибки обычно звучит так: «мы же уже обработали сообщение, почему оно пришло ещё раз?» Или наоборот: worker вызвал API, после чего задача пропала без результата. Цена обоих случаев одна — команда смешала два разных объекта: бизнес-операцию и delivery от broker. Первый живёт столько, сколько нужен отчёт или запись; второе живёт до ack/nack на конкретном канале. Пока между ними нет явного перехода, дубль кажется аварией, а ранний ack кажется оптимизацией.
В мае 2020 полезнее освоить небольшой автомат, чем рисовать сложную оркестрацию. У операции есть jobId, у каждой доставки — свой tag и признак redelivery, а worker делает четыре действия в фиксированном порядке: получить delivery, захватить состояние задачи, сохранить исход работы, сообщить broker результат обработки. Этот порядок не даёт «ровно один раз» на всём мире. Он даёт проверяемую at-least-once модель, в которой повтор не обязан повторять эффект одной уже завершённой задачи.
Delivery и задача отвечают на разные вопросы
Delivery отвечает на вопрос broker: кто сейчас несёт ответственность за эту копию сообщения? Для ручного подтверждения ответ заканчивается ack(tag) или отрицательным ответом. Если соединение consumer закрылось до ack, broker может вернуть непроверенную доставку и позже отдать её тому же или другому worker. Задача отвечает на вопрос приложения: что пользователь попросил, на какой попытке это находится и где результат. Её нельзя определить только по тому, есть ли сообщение в очереди.
Из этой разницы получается рабочее правило. Не используем delivery tag как jobId: tag привязан к каналу и меняется при следующей доставке. В payload кладём короткий стабильный идентификатор и версию контракта, а детали входа храним у владельца задачи либо подписываем настолько ясно, чтобы worker мог проверить их версию. Когда приходит redelivery, worker ищет тот же jobId, а не пытается угадать, была ли такая строка по совпадению времени или имени файла.
| Сущность | Живёт где | Меняется когда | Правильное применение |
|---|---|---|---|
jobId | в базе приложения, payload и журнале | не меняется между повторами | идемпотентность, статус, результат, поддержка |
| delivery tag | в канале broker/client | при каждом новом delivery | ровно один ack/nack конкретной доставки |
| attempt | в записи задачи или явно в повторном сообщении | при принятом решении о повторе | лимит retry и отчёт о причине |
| resultKey | в устойчивом storage/БД | один раз при успехе | доказательство, что effect уже готов до ack |
Префикс «at least once» относится к доставке, не к тому, что пользователь увидит два отчёта. Повтор возможен и после publisher-side неопределённости, и после consumer-side падения. Поэтому нельзя надеяться, что флаг redelivered всё решит: он полезен для журнала и приоритета проверки, но не заменяет запись состояния. Безопасное решение на уровне одной задачи — выбирать один стабильный эффект: например, файл всегда пишется в reports/{jobId}.csv, а строка результата хранит этот же ключ.
Ack ставим после устойчивого бизнес-результата
Ручной ack — это не запись в лог «worker начал работу». Это граница, после которой broker вправе удалить delivery. Если handler отправил ack сразу после получения, то сбой в SQL, файловом хранилище или внешнем API оставит ложный след: очередь считает задачу завершённой, приложение — нет. Если ack ставится после результата, возможен другой порядок: результат уже сохранён, а сеть закрылась. Следующее delivery должно обнаружить завершённую задачу и не делать effect второй раз.
Здесь важно назвать настоящий порядок устойчивости. Внутри приложения markSucceeded и resultKey должны быть сохранены одной согласованной операцией или в таком порядке, который можно восстановить после рестарта. Для внешнего вызова нужен собственный ключ: запись файла по детерминированному пути, уникальный ключ отправки или API-контракт, который принимает idempotency key. Нельзя сначала сделать необратимый эффект, а потом надеяться, что локальная таблица спасёт от дубля. Таблица лишь даёт worker решение при следующем delivery.
// Учебный обработчик одного delivery. transport API намеренно абстрактный.
async function handleDelivery(delivery, jobs, broker) {
const job = await jobs.findForUpdate(delivery.jobId);
if (job.state === 'succeeded' || job.state === 'quarantined') {
await broker.ack(delivery.tag);
return;
}
try {
await jobs.markRunning(job.id, delivery.attempt);
const resultKey = await writeReportOnce(job.id, job.payload);
await jobs.markSucceeded(job.id, resultKey);
await broker.ack(delivery.tag); // только после durable result + state
} catch (error) {
const nextAttempt = delivery.attempt + 1;
await jobs.markRetryOrQuarantine(job.id, nextAttempt, error.code);
await broker.nack(delivery.tag, { requeue: nextAttempt < 3 });
}
}
Это не интерфейс конкретной Node-библиотеки. Он намеренно показывает точки, которые нельзя переставлять: terminal state проверяется до действия; попытка фиксируется до повторного маршрута; success пишется до ack; карантин фиксируется до nack без requeue. Реальный API может называть методы иначе и по-разному задавать задержку. До внедрения надо сверить эти места с версией выбранного client и с тем, настроен ли у очереди путь dead-lettering.
Повтор классифицируем до того, как вернуть сообщение
Не всякая ошибка заслуживает retry. Timeout до получения ответа иногда временный, но он не доказывает, что внешняя операция не состоялась; здесь особенно нужен ключ идемпотентности. Ошибка в payload, неизвестная версия сообщения или нарушенное обязательное поле повтором не исправится. Ошибка локальной валидации должна сразу закончить processing как карантин или осознанная отмена, а не навсегда держать одно delivery на голове очереди.
| Класс | Пример симптома | Проверка перед решением | Следующее действие |
|---|---|---|---|
| Успех | resultKey уже сохранён | состояние terminal и эффект читается по jobId | ack; при дубле — только ack |
| Временная | короткий timeout до ответного байта | записаны attempt и причина; есть бюджет меньше лимита | перевести job в retry_wait, вернуть через ограниченный маршрут |
| Неопределённая внешняя | таймаут после отправки запроса | проверить внешний ключ или статус по jobId | не создавать второй эффект; повторять только через идемпотентный контракт |
| Постоянная | payload не проходит version/validation | причина воспроизводится на тех же данных | quarantine и nack/reject без requeue |
Задержка между попытками тоже часть контракта. В 2020 году для одного сервиса можно обойтись простой retry-очередью с TTL или расписанием, которое уже умеет конкретная библиотека; не обязательно строить платформу. Но у каждой задачи должны быть фиксированные максимум попыток, причина последнего перехода и время следующего допуска. Формула exponential backoff не спасает, если неизвестно, какая именно попытка уже была и кто вернёт сообщение из задержки.
Poison message — это работа, которая больше не должна горячо крутиться
В AMQP basic.reject с requeue=false не обозначает «успех». Оно заканчивает обработку данного delivery; при заранее настроенном dead-letter exchange broker направляет сообщение на отдельный маршрут, иначе оно может быть отброшено. Это требует явной проектной договорённости: куда попадёт карантин, кто посмотрит его, как сопоставить payload с записью задачи и как защитить чувствительные поля от попадания в журнал.
Не стоит строить логику на бесконечном nack(requeue=true). RabbitMQ прямо предупреждает, что consumer, который все время возвращает delivery, может создать затратный redelivery loop. Лимит попыток в записи задачи плюс отдельный результат quarantined разрывают этот цикл с понятной ценой: часть работы не завершена автоматически, зато полезные сообщения продолжают получать worker, а человек получает конкретную причину вместо бесконечного шума.
Короткий журнал заменяет догадку о порядке
Для первой версии не нужна отдельная observability-платформа. Достаточно, чтобы каждый переход писал один и тот же набор: event, jobId, attempt, признак redelivered, причина и при успехе resultKey. Тогда вопрос «был ли ack до результата?» проверяется порядком пяти строк, а не памятью того, кто разбирает сбой. Нельзя класть в такой журнал полный payload, токены или документ пользователя: для связи достаточно идентификатора и безопасной классификации ошибки.
{"event":"received","jobId":"export-42-2020-05","attempt":1,"redelivered":false}
{"event":"retry_scheduled","jobId":"export-42-2020-05","attempt":1,"reason":"upstream_timeout"}
{"event":"received","jobId":"export-42-2020-05","attempt":2,"redelivered":true}
{"event":"result_saved","jobId":"export-42-2020-05","resultKey":"reports/export-42-2020-05.csv"}
{"event":"ack_sent","jobId":"export-42-2020-05","attempt":2}
Это пример ожидаемого учебного журнала, не снятый log реального worker. Первая строка показывает начало первой попытки, вторая — сохранённое решение о повторе. Вторая доставка имеет тот же jobId, но уже другой delivery и redelivered: true. Только после строки result_saved возникает ack_sent. Если эти две строки поменялись местами, расследование закончено: worker подтверждает работу до того, как может доказать её эффект.
Маршрут проверки автомата
- Составить таблицу состояний для одного
jobId: queued, running, retry_wait, succeeded и quarantined; убрать неявное «вроде выполняется». - Взять один payload и показать два delivery с разными tags; убедиться, что оба ищут одну и ту же запись задачи.
- Смоделировать падение после
result_savedдо ack. Повтор должен только подтвердить новое delivery, а не записать второй результат. - Смоделировать временную ошибку до лимита и проверить, что attempt и причина сохранены до постановки retry.
- Смоделировать невалидный payload либо третий отказ: задача становится quarantined, а немедленный requeue не разрешён.
- На отдельном стенде выбранного broker проверить реальную семантику ack/nack, redelivery, policy DLX и задержки; учебный автомат этого не заменяет.
Граница модели
Этот материал не утверждает, что любая очередь предоставляет одинаковый delivery count, delayed retry или безопасный dead-letter путь. В нём использованы понятия AMQP 0-9-1 и документация RabbitMQ как проверяемый ориентир; конкретная конфигурация зависит от версии broker, типа очереди и client library. Внутри пакета выполнена только детерминированная фикстура переходов в памяти. Она не открывает соединение с RabbitMQ, не измеряет throughput, не запускает конкурирующих worker и не заменяет интеграционный тест проекта.
Проверяемые источники
- AMQP 0-9-1 specification: basic.ack и basic.reject — первичная спецификация: delivery tag адресует доставку, а basic.reject с requeue управляет возвратом или отказом от сообщения
- RabbitMQ: Consumer Acknowledgements and Publisher Confirms — официальное описание ручного ack, автоматического requeue не подтверждённой доставки и риска немедленного redelivery loop
- RabbitMQ: Reliability guide — официальная граница между подтверждением доставки и обработкой, а также необходимость идемпотентного consumer при повторной доставке
- RabbitMQ: Dead Letter Exchanges — официальное описание маршрутизации отклонённого сообщения в отдельный exchange и причин dead-lettering