diff --git a/editorial/agent-rewrites/274.json b/editorial/agent-rewrites/274.json index 11adafb..0dc1de2 100644 --- a/editorial/agent-rewrites/274.json +++ b/editorial/agent-rewrites/274.json @@ -1,7 +1,7 @@ { "index": 274, "slug": "editorial-2020-05-field-background-jobs", - "title": "Фоновая задача продублировала экспорт: порядок записи, ack и карантин ошибок", - "excerpt": "Разбираем два сбоя фонового worker: повторный внешний эффект после потери ack и бесконечный requeue для постоянной ошибки. Показываем порядок действий, журнал и критерий готовности.", - "contentHtml": "
Пользователь запускает экспорт и через несколько минут получает два одинаковых файла. В той же очереди задача с неизвестным типом отчёта появляется снова сразу после отказа worker. Первая ошибка создаёт лишний внешний эффект: приходится выяснять, какой файл считать правильным, и чистить дубликаты. Вторая забивает worker одинаковыми исключениями. Полезные задачи ждут, а журнал растёт быстрее, чем его успевает читать дежурный инженер.
\nТезис такой: доставка сообщения и выполнение бизнес-операции — разные факты. Broker может доставить одну бизнес-задачу повторно, если не увидел подтверждение. Worker обязан сделать повтор безопасным по устойчивому ключу. Он также обязан отличать временный сбой от постоянной ошибки входных данных. Для первой ошибки нужен идемпотентный результат и правильный порядок записи. Для второй — ограничение попыток и карантин, а не безусловный requeue.
Пусть запрос создаёт задачу с jobId = export-42-2020-05. Producer публикует сообщение, а worker получает delivery с отдельным deliveryTag. Первый идентификатор относится к бизнес-операции. Второй относится к конкретной доставке на канале. Новый deliveryTag не означает новый экспорт.
Worker начинает работу, сохраняет файл и переводит задачу в succeeded. Затем соединение с broker обрывается до ack. Для broker доставка осталась неподтверждённой. После восстановления consumer получает её снова. Если обработчик каждый раз вызывает экспорт, он создаёт второй файл. Если имя строится из текущего времени, по имени файла невозможно доказать, что это повтор.
Вторая задача содержит reportKind = unknown. Worker не может обработать её ни сейчас, ни через секунду: причина находится во входных данных. Если обработчик на любое исключение отвечает nack(requeue=true), broker снова выдаёт то же сообщение. Так возникает hot loop. Повтор помогает только тогда, когда новая попытка может изменить исход.
Подтверждение не является сигналом «код вошёл в функцию». Для ручного ack оно должно означать, что consumer выполнил работу, за которую берёт ответственность. Если ack отправить до записи результата, сбой после ack может потерять работу. Если записать результат, но не сделать повтор безопасным, сбой до ack создаст дубль. Поэтому сначала фиксируют бизнес-результат, затем подтверждают delivery.
\nОдного порядка недостаточно. Повтор должен найти тот же объект по jobId, прочитать terminal state и вернуть уже сохранённый результат. Внешний ключ результата тоже должен быть детерминированным: например, reports/export-42-2020-05.csv. Это защищает путь к файлу, но не любой внешний эффект. Email-провайдер или другой HTTP-сервис должен поддерживать собственный idempotency key либо получать вызов через отдельный надёжный протокол.
Ниже — учебный псевдокод. Он показывает границу ответственности, но не является готовым клиентом RabbitMQ, не задаёт транзакционный API базы и не сообщает показатели production-нагрузки.
\nasync function handle(delivery) {\n const { jobId, reportKind } = delivery.payload;\n const job = await jobs.lockById(jobId);\n\n if (job.state === 'succeeded') {\n await delivery.ack();\n return job.resultKey;\n }\n\n if (!isKnownReport(reportKind)) {\n await jobs.markQuarantined(jobId, {\n reason: 'unknown_report_kind',\n payloadVersion: delivery.payload.version,\n });\n await delivery.nack({ requeue: false });\n return;\n }\n\n const resultKey = `reports/${jobId}.csv`;\n await exportReport({ reportKind, resultKey });\n await jobs.markSucceeded(jobId, resultKey);\n await delivery.ack();\n}\nПроверка terminal state стоит до внешнего действия. Блокировка или другой механизм защиты строки должен охватывать чтение состояния и резервирование работы. В реальной системе нужно решить, что происходит при падении между записью файла и markSucceeded. Обычно результат пишут во временный объект, затем атомарно публикуют его по детерминированному ключу и сохраняют состояние. Конкретный storage может потребовать иной протокол.
Ветвь постоянной ошибки не вызывает экспорт и не отправляет сообщение обратно в основную очередь. Она сохраняет причину и версию входных данных до отрицательного подтверждения. Затем transport направляет сообщение в настроенный dead-letter route либо в другой согласованный карантин. Если такого маршрута нет, requeue=false не создаёт архив для разбора автоматически: сообщение может быть отброшено в зависимости от конфигурации.
| Симптом | Причина | Проверка | Действие |
|---|---|---|---|
| Для одного jobId появились два файла | Результат создаётся по случайному имени или проверка состояния стоит после эффекта | Сопоставить jobId, resultKey, порядок result_saved, succeeded и ack_sent | Читать terminal state до экспорта и сохранять результат по детерминированному ключу |
| Повторная доставка приходит после успешной записи | Соединение оборвалось до ack | Сверить журнал broker, delivery tag и запись задачи | Обработать повтор как тот же jobId, не запускать внешний эффект снова |
| attempt растёт с одной и той же причиной | Постоянную ошибку приняли за временную | Сравнить payload и код причины на соседних попытках | Сохранить причину, перевести задачу в quarantined, отключить requeue |
| После nack сообщение исчезло | Для очереди не настроен dead-letter маршрут или сообщение не должно храниться | Проверить policy, binding, routing key и журнал публикации | Настроить проверяемый карантин либо явно принять потерю как контракт |
| Задача долго стоит в очереди | Не сработал outbox, producer публикует не туда или worker не читает binding | Найти запись job, marker публикации и факт первой доставки | Проверить dispatcher и маршрут до handler; не менять handler без delivery |
Таблица задаёт порядок расследования. Сначала привяжите наблюдение к одному jobId. Не начинайте с увеличения timeout или числа retry. Эти изменения не исправят случайное имя результата, неизвестный reportKind или неверный binding.
Ручное подтверждение помогает выбрать границу ответственности, но не превращает распределённую систему в exactly-once механизм. Между записью результата и ack остаётся окно сбоя. Между отправкой ack и получением его broker тоже есть сеть. Система с повторной доставкой обычно даёт at-least-once обработку: сообщение могут обработать снова, поэтому handler должен быть идемпотентным.
\nПризнак redelivered полезен для диагностики, но он не заменяет бизнес-ключ. При повторной публикации похожее сообщение может выглядеть как новая доставка. Бизнес-правило должно опираться на jobId, уникальный ключ результата и состояние операции. Если worker выполняют несколько экземпляров, lock и уникальное ограничение должны защищать одну и ту же область.
Для временной ошибки используйте ограниченный retry с задержкой. Временной может быть недоступность хранилища, краткий сетевой отказ или ещё не созданная зависимая запись. Сохраняйте attempt, код причины и время следующей попытки. После лимита не отправляйте сообщение в бесконечный цикл. Оставьте его в карантине, чтобы инженер мог исправить данные или код и принять решение о повторном запуске.
jobId. Соберите запись задачи, payload, outbox, delivery и журнал worker в одной временной шкале.result_saved, terminal state и ack_sent. Не запускайте экспорт повторно, пока не проверили сохранённый результат.Описанный алгоритм не говорит, какая база или очередь нужна проекту. Row lock защищает только выбранную запись и не заменяет уникальное ограничение во внешнем хранилище. Состояние в базе и файл в object storage не откатываются одной транзакцией. Если процесс остановился после загрузки файла, нужна проверка существующего ключа и политика очистки незавершённых объектов.
\nЛимит retry нельзя выбрать одинаковым для всех задач. Слишком малый лимит превращает краткий сбой в карантин. Слишком большой растягивает задержку и маскирует постоянную ошибку. Backoff и лимит должны учитывать стоимость работы, срок жизни входных данных и допустимую задержку пользователя.
\nУчебный код не доказывает поведение конкретной версии broker, client library или policy. Проверяйте ack, redelivery, nack, dead-lettering и восстановление соединения на той конфигурации, которую действительно запускает сервис. Не называйте результат готовым, если проверили только in-memory переходы.
\nРешение готово, когда два независимых сценария проходят на выбранном окружении. После сбоя между сохранением результата и ack повторная доставка того же jobId не создаёт второй внешний эффект и возвращает один проверяемый resultKey. После постоянной ошибки payload задача достигает quarantined не позднее заданного лимита, причина сохраняется, а сообщение не возвращается в основную очередь. В журнале можно восстановить порядок событий без догадок.
Симптом в учебном разборе такой: пользователь запрашивает экспорт, а через несколько минут видит два одинаковых файла. Одновременно другая заявка с неизвестным типом отчёта снова и снова появляется у worker, занимая очередь. Цена двойная. Первый сбой создаёт лишний внешний эффект и спор, какая копия верная; второй забирает время worker и скрывает полезные задачи под повторяющейся ошибкой. Фраза «очередь доставила дважды» описывает факт, но ещё не называет место, где принято неверное решение.
\nРазберу не production-инцидент, а анонимизированную in-memory фикстуру мая 2020 года. В ней нет реального broker, файлового хранилища, user data или измеренной нагрузки. Зато есть две детерминированные цепочки с одним jobId каждая: временный отказ между сохранением результата и ack, а также невалидный payload после лимита попыток. Цель — показать практическое расследование: симптом → причина → проверка → действие, а не рассказать историю успеха постфактум.
Первый экспорт имеет 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}\n{"event":"validation_failed","jobId":"export-43-2020-05","reason":"unknown_report_kind"}\n{"event":"quarantined","jobId":"export-43-2020-05","queue":"jobs.quarantine"}\n{"event":"nack_sent","jobId":"export-43-2020-05","requeue":false}\nСтроки выше — синтетический журнал фикстуры, не вывод запущенного 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 повторно requeue | validation_failed повторяется на равных входных данных | пометить quarantined и направить в DLX/разбор |
| queued долго без received | outbox не опубликован либо worker не читает маршрут | есть job/outbox, но нет publish marker и journal delivery | проверить dispatcher, binding и конкретный broker client |
Эта таблица не заменяет доступ к очереди. Она задаёт порядок вопросов до изменения кода. Если уже есть resultKey, не нужно стартовать новый экспорт «на всякий случай». Если причина unknown_report_kind воспроизводится из сохранённой версии payload, не нужно увеличивать retry. Если у задачи нет received, бесполезно рассматривать handler: сперва ищем outbox, публикацию и маршрут. Каждый шаг привязывает действие к одному наблюдаемому факту.
Плохой вариант handler выглядит почти естественно: он берёт сообщение, сразу начинает генерацию, формирует имя из текущего времени, отправляет ack и только затем пытается отметить успех. В нём две точки неопределённости. Во-первых, повтору нечем доказать, что прежняя генерация уже завершилась: название файла другое, а состояние ещё не terminal. Во-вторых, ack может добраться до broker раньше записи статуса. При сбое получаем либо потерянную работу, либо повтор без защиты.
\nИсправление на уровне одной задачи не требует общего дедупликатора. 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) всегда | горячий цикл на невалидном payload | classify → retry_wait или quarantined → nack по решению | attempt ограничен, причина остаётся рядом с jobId |
| искать ошибку по времени | непонятно, к какой попытке относится строка | писать jobId, attempt, event, reason | одна цепочка читается без догадки о совпадении |
Здесь нет обещания, что SQL-блокировка сделает worker глобально одиночным. Она лишь защищает запись задачи в границе выбранной базы. Конкурирующие worker, внешнее хранилище и сеть всё равно требуют проверяемого контракта. Поэтому практический критерий короче: каждый новый delivery для уже succeeded обязан завершиться без нового результата. Если это нельзя проверить, слово «идемпотентность» в код-ревью пока ничего не означает.
Повтор нужен, когда новое время может изменить исход: зависимость была недоступна, лимит снят, ожидаемая запись ещё не появилась. Но unknown_report_kind не зависит от времени. Для такого случая worker должен назвать ошибку постоянной, сохранить её в записи и завершить автоматический путь. В AMQP отрицательный ответ без requeue может направить сообщение в DLX, если проект это настроил. Отдельный маршрут делает ошибку предметом разбора, а не бесконечным consumer workload.
Карантин не стоит использовать как корзину для всех исключений. Сначала сохраняем тип причины и версию payload: это отделяет неисправимые данные от дефекта worker после обновления. Затем владелец решает: поправить данные и переиздать новую задачу с новым или тем же бизнес-ключом, починить consumer и вручную вернуть сообщение через контролируемый маршрут, либо отменить операцию. Автоматическое чтение карантина обратно в основную очередь без исправления причины снова создаёт тот же loop, только с более длинным названием.
\nВ модуле ревизий есть команда 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}\n{"event":"retry_scheduled","jobId":"export-42-2020-05","attempt":1,"reason":"upstream_timeout"}\n{"event":"received","jobId":"export-42-2020-05","attempt":2,"redelivered":true}\n{"event":"result_saved","jobId":"export-42-2020-05","resultKey":"reports/export-42-2020-05.csv"}\n{"event":"ack_sent","jobId":"export-42-2020-05","attempt":2}\nФикстура не даёт ложной уверенности в broker. У неё нет TCP-соединения, реального delivery tag, политики DLX, нескольких consumer или диска. Но она отделяет две логические проверки, которые можно выполнить без инфраструктуры: для retry есть новая попытка с тем же jobId, а успешный путь пишет result до ack; poison-путь обрывает requeue на известном пределе. После выбора библиотеки эту же пару сценариев нужно повторить на интеграционном стенде и сравнить реальные журналы с ожидаемыми переходами.
\nВсе идентификаторы, причины и строки журнала здесь придуманы для проверки переходов. Нет реального файла, заказчика, очереди, RabbitMQ policy, production-config или browser-действия. Тексты опираются на спецификацию AMQP и официальную документацию RabbitMQ, чтобы не выдумывать смысл ack, reject и redelivery, но не выдают современную документацию за снимок конкретной инфраструктуры мая 2020 года. Автор этого периода умеет провести узкое backend/delivery расследование и оставить route для стенда; он ещё не заявляет готовую платформу наблюдаемости или сложную оркестрацию.
\n