DarkRiDDeR14 мин

Разбор: почему экспорт отчёта повторился и как остановить poison message

BackendОчередиРазбор

Симптом в учебном разборе такой: пользователь запрашивает экспорт, а через несколько минут видит два одинаковых файла. Одновременно другая заявка с неизвестным типом отчёта снова и снова появляется у worker, занимая очередь. Цена двойная. Первый сбой создаёт лишний внешний эффект и спор, какая копия верная; второй забирает время worker и скрывает полезные задачи под повторяющейся ошибкой. Фраза «очередь доставила дважды» описывает факт, но ещё не называет место, где принято неверное решение.

Разберу не production-инцидент, а анонимизированную in-memory фикстуру мая 2020 года. В ней нет реального broker, файлового хранилища, user data или измеренной нагрузки. Зато есть две детерминированные цепочки с одним jobId каждая: временный отказ между сохранением результата и ack, а также невалидный payload после лимита попыток. Цель — показать практическое расследование: симптом → причина → проверка → действие, а не рассказать историю успеха постфактум.

Сначала отделяем факт результата от факта delivery

Первый экспорт имеет jobId = export-42-2020-05. Worker получил delivery, записал CSV по устойчивому ключу и пометил задачу как succeeded. Затем соединение до broker оборвалось до ack. У broker остаётся непроверенное delivery, поэтому следующий worker получает ту же бизнес-задачу повторно. Если handler относится к любому received message как к новому, он снова вызывает экспорт и пишет второй файл с новым случайным именем. Это не исправляется большим timeout: проблема в том, что idempotency check находится после эффекта или отсутствует.

Вторая заявка — export-43-2020-05 — содержит неизвестный reportKind. Worker ловит ошибку, делает nack(requeue=true) и тут же получает ту же доставку снова. Никакая пауза не сделает неизвестный тип валидным. Пока задача не имеет состояния quarantined и ограничителя попыток, очередь по сути работает как генератор одинаковых ошибок. Здесь цена уже операционная: журнал растёт, полезная работа ждёт, а владелец данных не получает короткий список того, что нужно исправить.

{"event":"received","jobId":"export-43-2020-05","attempt":3,"redelivered":true}
{"event":"validation_failed","jobId":"export-43-2020-05","reason":"unknown_report_kind"}
{"event":"quarantined","jobId":"export-43-2020-05","queue":"jobs.quarantine"}
{"event":"nack_sent","jobId":"export-43-2020-05","requeue":false}

Строки выше — синтетический журнал фикстуры, не вывод запущенного RabbitMQ consumer. Они важны именно порядком. У poison-задачи третья попытка ещё фиксирует вход и причину, затем приложение сохраняет quarantined, и только после этого выбранному transport посылается отрицательный ответ без requeue. В реальном AMQP дальнейшая судьба зависит от настроенного dead-letter exchange: без маршрута сообщение может быть отброшено. Поэтому «карантин» обязан существовать не только как слово в коде, но и как проверяемая конфигурация выбранного окружения.

Карта расследования: какой факт исключает какую гипотезу
НаблюдениеПричина, которую проверяемМинимальное доказательствоДействие
Есть resultKey, но нет ack_sentсбой произошёл в узком окне после результатаjob.state = succeeded раньше следующего receivedпри повторе не создавать результат, подтвердить новое delivery
Два файла для одного jobIdвнешний эффект не имеет стабильного ключа либо check сделан поздносопоставить имена файлов и порядок journalпуть результата построить из jobId, terminal state читать до export
attempt растёт, причина одна и та жепостоянный payload повторно requeuevalidation_failed повторяется на равных входных данныхпометить quarantined и направить в DLX/разбор
queued долго без receivedoutbox не опубликован либо worker не читает маршрутесть job/outbox, но нет publish marker и journal deliveryпроверить dispatcher, binding и конкретный broker client

Эта таблица не заменяет доступ к очереди. Она задаёт порядок вопросов до изменения кода. Если уже есть resultKey, не нужно стартовать новый экспорт «на всякий случай». Если причина unknown_report_kind воспроизводится из сохранённой версии payload, не нужно увеличивать retry. Если у задачи нет received, бесполезно рассматривать handler: сперва ищем outbox, публикацию и маршрут. Каждый шаг привязывает действие к одному наблюдаемому факту.

Вертикальная схема диагностики фоновой задачи: по jobId проверяют запись задачи, outbox и сообщение, затем журнал worker, сохранённый resultKey и при постоянной ошибке карантинную очередь
Разбор идёт не по названию компонента, а по пути одного jobId. Это уменьшает риск перезапустить эффект, когда проблема находится до worker или после результата.

Причина первого дубля: случайное имя и ack не на той стороне

Плохой вариант handler выглядит почти естественно: он берёт сообщение, сразу начинает генерацию, формирует имя из текущего времени, отправляет ack и только затем пытается отметить успех. В нём две точки неопределённости. Во-первых, повтору нечем доказать, что прежняя генерация уже завершилась: название файла другое, а состояние ещё не terminal. Во-вторых, ack может добраться до broker раньше записи статуса. При сбое получаем либо потерянную работу, либо повтор без защиты.

Исправление на уровне одной задачи не требует общего дедупликатора. Worker сначала читает строку по jobId с блокировкой, проверяет terminal states и резервирует попытку. Результат записывается по детерминированному ключу reports/{jobId}.csv. После записи в этом же бизнес-шаге сохраняются resultKey и succeeded. Если после этого broker повторно доставит сообщение, handler видит terminal state, не пишет файл заново и только завершает текущий delivery. Для email или внешнего API нужно отдельно убедиться, что принимающая сторона поддерживает такой ключ; путь файла не решает чужой side effect.

Изменение порядка для повторной доставки
Старая последовательностьРискНовая последовательностьПроверяемый результат
receive → generate random file → ack → save stateдубль или потеря при падении между шагамиreceive → read job → save deterministic result + succeeded → ackповтор видит succeeded и не создаёт второй файл
catch → nack(requeue=true) всегдагорячий цикл на невалидном payloadclassify → retry_wait или quarantined → nack по решениюattempt ограничен, причина остаётся рядом с jobId
искать ошибку по временинепонятно, к какой попытке относится строкаписать jobId, attempt, event, reasonодна цепочка читается без догадки о совпадении

Здесь нет обещания, что SQL-блокировка сделает worker глобально одиночным. Она лишь защищает запись задачи в границе выбранной базы. Конкурирующие worker, внешнее хранилище и сеть всё равно требуют проверяемого контракта. Поэтому практический критерий короче: каждый новый delivery для уже succeeded обязан завершиться без нового результата. Если это нельзя проверить, слово «идемпотентность» в код-ревью пока ничего не означает.

Причина hot loop: постоянную ошибку приняли за временную

Повтор нужен, когда новое время может изменить исход: зависимость была недоступна, лимит снят, ожидаемая запись ещё не появилась. Но unknown_report_kind не зависит от времени. Для такого случая worker должен назвать ошибку постоянной, сохранить её в записи и завершить автоматический путь. В AMQP отрицательный ответ без requeue может направить сообщение в DLX, если проект это настроил. Отдельный маршрут делает ошибку предметом разбора, а не бесконечным consumer workload.

Карантин не стоит использовать как корзину для всех исключений. Сначала сохраняем тип причины и версию payload: это отделяет неисправимые данные от дефекта worker после обновления. Затем владелец решает: поправить данные и переиздать новую задачу с новым или тем же бизнес-ключом, починить consumer и вручную вернуть сообщение через контролируемый маршрут, либо отменить операцию. Автоматическое чтение карантина обратно в основную очередь без исправления причины снова создаёт тот же loop, только с более длинным названием.

Фикстура как маленький регрессионный контракт

В модуле ревизий есть команда node scripts/upgrade-2020-05.mjs --verify-fixture. Она не эмулирует AMQP frames. Она детерминированно создаёт две записи в памяти: первая переживает transient ошибку, получает повторную доставку, сохраняет result и ack; вторая после третьего неуспеха переходит в quarantined. Вывод — JSON-журнал и три булевых условия. Такой тест полезен тем, что не позволяет незаметно переставить result_saved и ack_sent в учебном алгоритме.

{"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}

Фикстура не даёт ложной уверенности в broker. У неё нет TCP-соединения, реального delivery tag, политики DLX, нескольких consumer или диска. Но она отделяет две логические проверки, которые можно выполнить без инфраструктуры: для retry есть новая попытка с тем же jobId, а успешный путь пишет result до ack; poison-путь обрывает requeue на известном пределе. После выбора библиотеки эту же пару сценариев нужно повторить на интеграционном стенде и сравнить реальные журналы с ожидаемыми переходами.

Маршрут разбора перед исправлением

  1. Взять один конкретный jobId и собрать рядом запись задачи, outbox, журнал worker, ключ результата и информацию о current attempt.
  2. Проверить, в каком порядке появились result_saved, succeeded и ack_sent; не делать новый экспорт, пока это не ясно.
  3. Для повтора сравнить jobId, а не delivery tag: новый tag не означает новую бизнес-операцию.
  4. Классифицировать последнюю ошибку как временную, неопределённую внешнюю или постоянную; записать основание рядом с attempt.
  5. Для постоянной ошибки остановить requeue, перевести job в quarantined и проверить, что выбранная DLX/карантинная поверхность действительно принимает сообщение.
  6. После изменения прогнать in-memory фикстуру, затем отдельный broker-интеграционный сценарий с падением до ack; в этом пакете выполнен только первый шаг.

Граница полевого разбора

Все идентификаторы, причины и строки журнала здесь придуманы для проверки переходов. Нет реального файла, заказчика, очереди, RabbitMQ policy, production-config или browser-действия. Тексты опираются на спецификацию AMQP и официальную документацию RabbitMQ, чтобы не выдумывать смысл ack, reject и redelivery, но не выдают современную документацию за снимок конкретной инфраструктуры мая 2020 года. Автор этого периода умеет провести узкое backend/delivery расследование и оставить route для стенда; он ещё не заявляет готовую платформу наблюдаемости или сложную оркестрацию.

Проверяемые источники

  • 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