diff --git a/editorial/agent-rewrites/275.json b/editorial/agent-rewrites/275.json index 1e4061b..4a926c6 100644 --- a/editorial/agent-rewrites/275.json +++ b/editorial/agent-rewrites/275.json @@ -3,5 +3,5 @@ "slug": "editorial-2020-05-mechanism-background-jobs", "title": "Фоновая задача без дублей: delivery, ack и карантин", "excerpt": "Как отделить бизнес-состояние фоновой задачи от доставки сообщения, поставить ack после результата и остановить бесконечный retry для неисправимых данных.", - "contentHtml": "

Пользователь запускает экспорт отчёта и получает два одинаковых файла. В другом случае worker вызывает внешний API, после чего сообщение исчезает, а результата нет. Оба симптома появляются на одной границе: приложение путает бизнес-задачу с отдельной доставкой сообщения. Цена ошибки — повторный платный вызов, дублирующий файл или письмо, потерянная работа и очередь, забитая одной неисправимой записью.

\n

Рабочая модель разделяет эти объекты. Бизнес-задача имеет стабильный jobId, состояние и ключ результата. Broker доставляет сообщение с собственным delivery tag. Worker проверяет состояние, выполняет эффект с устойчивым ключом, сохраняет результат и только потом подтверждает конкретную доставку через ack. Если связь оборвётся до подтверждения, broker может доставить сообщение снова. Повтор не должен создавать новый эффект для уже завершённого jobId.

\n

Delivery и задача отвечают на разные вопросы

\n

Delivery отвечает на вопрос broker: кто сейчас отвечает за эту копию сообщения? Его жизненный цикл заканчивается на ack, reject или закрытии канала. При закрытии канала неподтверждённая доставка может вернуться в очередь.

\n

Задача отвечает на вопрос приложения: что попросил пользователь, на какой попытке находится операция и где лежит результат. Её состояние должно жить в базе или другом устойчивом хранилище. Нельзя использовать delivery tag как идентификатор задачи. Tag относится к каналу и может измениться при следующей доставке той же бизнес-операции.

\n
Идентификаторы фоновой обработки
ОбъектГде живётКогда меняетсяДля чего нужен
jobIdзапись задачи, payload, журналне меняется между повторамисостояние, дедупликация и результат
delivery tagканал consumerпри новой доставкеточный ack или nack
attemptзапись задачи или retry-сообщениепри разрешённой новой попыткелимит повторов и диагностика
resultKeyбаза и хранилище результатаодин раз при успехедоказательство готового эффекта
\n

Из этого разделения следует ограничение: модель даёт at-least-once delivery, а не глобальный «ровно один раз». Сообщение может прийти повторно. Поэтому эффект должен быть идемпотентным в границе задачи. Для отчёта это может быть путь reports/{jobId}.csv и уникальная запись результата. Для внешнего API нужен его собственный idempotency key. Локальная таблица не отменяет уже отправленный запрос в чужую систему.

\n
\"Схема
Повтор относится к delivery, а состояние — к jobId. Сначала приложение сохраняет решение, затем broker получает подтверждение или отказ.
\n

Минимальный контракт задачи

\n

До публикации сообщения приложение создаёт запись задачи и outbox-событие в одной транзакции. Outbox хранит намерение опубликовать сообщение, пока dispatcher не получит подтверждение от выбранного broker-клиента. Такой порядок закрывает отдельную дыру: задача уже видна пользователю, но процесс публикации ещё не завершён.

\n
// Учебный псевдокод: это не готовый API RabbitMQ.\\nasync function requestExport(input, db) {\\n  const jobId = stableId(input.accountId, input.period);\\n  await db.transaction(async (tx) => {\\n    await tx.insertJob({ id: jobId, state: 'queued', attempt: 0, resultKey: null });\\n    await tx.insertOutbox({ type: 'report.export.requested', jobId, publishedAt: null });\\n  });\\n  return { accepted: true, jobId };\\n}
\n

Пример учебный. Он показывает контракт, а не измеренную производительность и не конкретную библиотеку. Функция принимает намерение и возвращает jobId. Она не держит HTTP-соединение до окончания экспорта. Dispatcher отдельно публикует событие и отмечает publishedAt после подтверждения своего клиентского API.

\n

Ack ставим после устойчивого результата

\n

Ранний ack сообщает broker, что доставка обработана. Если отправить его сразу после чтения сообщения, а затем получить ошибку базы, файлового хранилища или внешнего API, broker удалит delivery, хотя бизнес-результата нет. Это путь к потере работы.

\n

Поздний ack оставляет другое окно. Worker может сохранить результат, а соединение оборвётся до подтверждения. Broker доставит сообщение повторно. Второй worker должен прочитать terminal state и завершить только новое delivery. Он не должен повторять экспорт.

\n
// Учебный обработчик одного delivery.\\nasync function handleDelivery(delivery, jobs, broker) {\\n  const job = await jobs.findForUpdate(delivery.jobId);\\n  if (job.state === 'succeeded' || job.state === 'quarantined') {\\n    await broker.ack(delivery.tag);\\n    return;\\n  }\\n  try {\\n    await jobs.markRunning(job.id, delivery.attempt);\\n    const resultKey = await writeReportOnce(job.id, job.payload);\\n    await jobs.markSucceeded(job.id, resultKey);\\n    await broker.ack(delivery.tag);\\n  } catch (error) {\\n    const nextAttempt = delivery.attempt + 1;\\n    await jobs.markRetryOrQuarantine(job.id, nextAttempt, error.code);\\n    await broker.nack(delivery.tag, { requeue: nextAttempt < 3 });\\n  }\\n}
\n

Порядок в примере — часть контракта. Состояние succeeded проверяется до эффекта. Результат получает детерминированный ключ. Статус успеха сохраняется до ack. Для внешнего вызова нужна такая же защита на стороне API: ключ операции, уникальное ограничение или запрос статуса по прежнему ключу.

\n

Симптом → причина → проверка → действие

\n
Карта диагностики одного jobId
СимптомПричинаПроверкаДействие
Есть resultKey, но нет ackСвязь оборвалась после результатаСравнить порядок result_saved и ack_sentПри повторе прочитать terminal state и подтвердить только delivery
Два файла для одного jobIdСлучайное имя или поздняя проверкаСопоставить имена файлов с журналом workerИспользовать resultKey и проверять статус до эффекта
attempt растёт с одной причинойПостоянную ошибку отправляют в requeueПовторить validation на сохранённом payloadПеревести задачу в quarantined и прекратить requeue
Задача долго queuedНе сработал outbox или маршрутПроверить outbox, marker публикации и bindingИсправить dispatcher или маршрут, не менять handler вслепую
\n

Проверку ведут по одному jobId. В журнале достаточно событий received, result_saved, retry_scheduled, quarantined и ack_sent. Рядом пишут попытку, признак redelivery и безопасную причину. Полный payload, токены и пользовательские документы в журнал не кладут.

\n

Retry нужен не для любой ошибки

\n

Повтор оправдан, если новое время может изменить исход: зависимость временно недоступна, сработал сетевой timeout до ответа или ожидаемая запись ещё не появилась. Но timeout после отправки запроса не доказывает, что внешний эффект не состоялся. Такой вызов повторяют только с ключом идемпотентности или после проверки статуса операции.

\n

Невалидный JSON, неизвестная версия события и отсутствующее обязательное поле повтором не исправятся. Бесконечный nack(requeue=true) создаёт горячий redelivery loop. Он занимает worker и прячет полезные сообщения за одной постоянной ошибкой.

\n

Для retry задают максимальное число попыток, причину последнего перехода и время следующего допуска. Задержка может использовать retry-очередь с TTL или другой механизм выбранного клиента. Важно, чтобы worker знал текущую попытку и не возвращал неисправимую запись в основной маршрут без изменения причины.

\n

Карантин для poison message

\n

Poison message — сообщение, которое текущий consumer не может обработать автоматически. Worker сначала сохраняет причину и состояние quarantined, затем отклоняет delivery без requeue. При настроенном dead-letter exchange broker направит сообщение на отдельный маршрут. Без такой конфигурации оно может быть отброшено, поэтому карантин должен быть проверяемой частью инфраструктуры, а не только словом в коде.

\n

Карантин не означает успех. Он означает, что автоматический путь остановился с понятной причиной. Владелец может исправить payload и переиздать задачу, обновить consumer или отменить операцию. Автоматически читать карантин обратно в основную очередь без исправления причины нельзя: loop вернётся.

\n

Порядок действий

\n
  1. Определить стабильный jobId и записывать его в задачу, событие, результат и журнал.
  2. Разделить состояния queued, running, retry_wait, succeeded и quarantined.
  3. Проверить согласованное создание outbox и задачи, а также публикацию после подтверждения broker-клиента.
  4. Поставить проверку terminal state до необратимого эффекта.
  5. Сохранить результат и resultKey до ack конкретного delivery.
  6. Разделить временные, неопределённые и постоянные ошибки; для каждой задать проверку и лимит.
  7. Для постоянной ошибки записать quarantined, отправить отказ без requeue и проверить dead-letter маршрут.
  8. Прогнать повторную доставку после сохранённого результата и убедиться, что второй эффект не создаётся.
\n

Ограничения

\n

Эта схема не делает систему ровно-однократной. Два worker могут одновременно увидеть незахваченную задачу, если хранилище не даёт блокировку или уникальное ограничение. Внешний сервис может принять запрос и не вернуть ответ. Broker может иметь другую семантику подтверждений. Поэтому порядок нужно сверить с версией клиента, типом очереди и реальной политикой dead-lettering.

\n

Псевдокод выше не открывает соединение с RabbitMQ, не измеряет throughput и не является production-тестом. Он ограничен учебной иллюстрацией переходов. Интеграционная проверка должна использовать выбранный broker, несколько worker, падение до ack, повтор после результата и невалидный payload после лимита.

\n

Критерий готовности

\n

Решение готово, когда для одного заранее известного jobId журнал показывает устойчивый результат до первого ack, повторную доставку с новым tag и отсутствие второго эффекта. Для временной ошибки видны ограниченные попытки и следующий допуск. Для невалидного payload видны причина, состояние quarantined и отсутствие немедленного requeue. Эти свойства должны воспроизводиться на интеграционном стенде выбранного broker, а не только в unit-тесте.

\n

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

\n" + "contentHtml": "

Пользователь запускает экспорт отчёта и получает два одинаковых файла. В другом случае worker вызывает внешний API, после чего сообщение исчезает, а результата нет. Оба симптома появляются на одной границе: приложение путает бизнес-задачу с отдельной доставкой сообщения. Цена ошибки — повторный платный вызов, дублирующий файл или письмо, потерянная работа и очередь, забитая одной неисправимой записью.

\n

Рабочая модель разделяет эти объекты. Бизнес-задача имеет стабильный jobId, состояние и ключ результата. Broker доставляет сообщение с собственным delivery tag. Worker проверяет состояние, выполняет эффект с устойчивым ключом, сохраняет результат и только потом подтверждает конкретную доставку через ack. Если связь оборвётся до подтверждения, broker может доставить сообщение снова. Повтор не должен создавать новый эффект для уже завершённого jobId.

\n

Delivery и задача отвечают на разные вопросы

\n

Delivery отвечает на вопрос broker: кто сейчас отвечает за эту копию сообщения? В AMQP 0-9-1 delivery tag уникален только внутри канала и передаётся в ack, reject или nack. Поэтому подтверждать доставку нужно на том же канале, на котором она пришла. При закрытии канала неподтверждённая доставка может вернуться в очередь.

\n

Задача отвечает на вопрос приложения: что попросил пользователь, на какой попытке находится операция и где лежит результат. Её состояние должно жить в базе или другом устойчивом хранилище. Нельзя использовать delivery tag как идентификатор задачи. При новой доставке tag относится к тому же каналу, но идентифицирует уже новую delivery; признак повторной доставки передаётся отдельно.

\n
Идентификаторы фоновой обработки
ОбъектГде живётКогда меняетсяДля чего нужен
jobIdзапись задачи, payload, журналне меняется между повторамисостояние, дедупликация и результат
delivery tagканал consumerдля каждой deliveryточный ack или nack на том же канале
attemptзапись задачи или retry-сообщениепри разрешённой новой попыткелимит повторов и диагностика
resultKeyбаза и хранилище результатаодин раз при успехедоказательство готового эффекта
\n

Из этого разделения следует ограничение: модель даёт at-least-once delivery, а не глобальный «ровно один раз». Сообщение может прийти повторно. Поэтому эффект должен быть идемпотентным в границе задачи. Для отчёта это может быть путь reports/{jobId}.csv и уникальная запись результата. Для внешнего API нужен его собственный idempotency key. Локальная таблица не отменяет уже отправленный запрос в чужую систему.

\n
\"Схема
Повтор относится к delivery, а состояние — к jobId. Сначала приложение сохраняет решение, затем broker получает подтверждение или отказ.
\n

Минимальный контракт задачи

\n

До публикации сообщения приложение создаёт запись задачи и outbox-событие в одной транзакции. Outbox хранит намерение опубликовать сообщение, пока dispatcher не получит подтверждение от broker-клиента. Такое подтверждение говорит, что broker принял публикацию по своему контракту; оно не доказывает, что consumer уже выполнил бизнес-эффект. Задача может быть видна пользователю, хотя публикация ещё не завершена.

\n
// Учебный псевдокод: это не готовый API RabbitMQ.\nasync function requestExport(input, db) {\n  // В реальном проекте уникальность задаёт выбранная бизнес-граница.\n  const jobId = stableId(input.accountId, input.idempotencyKey);\n  await db.transaction(async (tx) => {\n    await tx.insertJob({\n      id: jobId,\n      state: 'queued',\n      attempt: 0,\n      resultKey: null,\n      payload: input.payload,\n    });\n    await tx.insertOutbox({\n      type: 'report.export.requested',\n      jobId,\n      publishedAt: null,\n    });\n  });\n  return { accepted: true, jobId };\n}
\n

Пример учебный. Он показывает контракт, а не измеренную производительность и не конкретную библиотеку. idempotencyKey должен иметь понятную область уникальности: иначе повтор запроса и новая операция за тот же период можно случайно слить. Функция принимает намерение и возвращает jobId. Она не держит HTTP-соединение до окончания экспорта. Dispatcher отдельно публикует событие и отмечает publishedAt после подтверждения своего клиентского API.

\n

Ack ставим после устойчивого результата

\n

Ранний ack сообщает broker, что доставка обработана. Если отправить его сразу после чтения сообщения, а затем получить ошибку базы, файлового хранилища или внешнего API, broker удалит delivery, хотя бизнес-результата нет. Это путь к потере работы.

\n

Поздний ack оставляет другое окно. Worker может сохранить результат, а соединение оборвётся до подтверждения. Broker доставит сообщение повторно. Второй worker должен прочитать terminal state и подтвердить только новое delivery. Он не должен повторять экспорт. Чтобы два worker не выполняли эффект одновременно, нужна атомарная блокировка или lease в хранилище; одного чтения состояния недостаточно.

\n
// Учебный обработчик одного delivery.\n// jobs.claim атомарно захватывает job или возвращает terminal state.\nasync function handleDelivery(delivery, jobs, broker) {\n  const job = await jobs.claim(delivery.jobId);\n  if (job.state === 'succeeded' || job.state === 'quarantined' || job.state === 'retry_wait') {\n    await broker.ack(delivery.tag);\n    return;\n  }\n  try {\n    const resultKey = await writeReportOnce(job.id, job.payload);\n    await jobs.markSucceeded(job.id, resultKey);\n    await broker.ack(delivery.tag);\n  } catch (error) {\n    // attempt хранится в задаче, а не в AMQP delivery.\n    const nextAttempt = job.attempt + 1;\n    if (isTransient(error) && nextAttempt < 3) {\n      // Состояние и retry-intent записываются одной транзакцией.\n      await jobs.markRetryWaitWithOutbox(job.id, nextAttempt, error.code);\n    } else {\n      await jobs.markQuarantined(job.id, nextAttempt, error.code);\n    }\n    // Задержанный retry публикуется отдельным dispatcher; не requeue-им delivery сразу.\n    await broker.nack(delivery.tag, { requeue: false });\n  }\n}
\n

Порядок в примере — часть контракта. claim должен атомарно исключать второй эффект или выдавать lease с понятным истечением; это проектная операция хранилища, а не гарантия RabbitMQ. Terminal state и retry_wait проверяются до эффекта. Результат получает детерминированный ключ. Статус успеха сохраняется до ack. Если nack потеряется после записи retry-intent, повторная delivery увидит retry_wait и подтвердится без нового эффекта. Для внешнего вызова нужна такая же защита на стороне API: ключ операции, уникальное ограничение или запрос статуса по прежнему ключу.

\n

Симптом → причина → проверка → действие

\n
Карта диагностики одного jobId
СимптомПричинаПроверкаДействие
Есть resultKey, но нет ackСвязь оборвалась после результатаСравнить порядок result_saved и ack_sentПри повторе прочитать terminal state и подтвердить только delivery
Два файла для одного jobIdСлучайное имя, гонка или поздняя проверкаСопоставить имена файлов с журналом worker и leaseИспользовать resultKey, уникальное ограничение и claim до эффекта
attempt растёт с одной причинойПостоянную ошибку отправляют в requeueПовторить validation на сохранённом payloadПеревести задачу в quarantined и прекратить requeue
Задача долго queuedНе сработал outbox или маршрутПроверить outbox, marker публикации и bindingИсправить dispatcher или маршрут, не менять handler вслепую
\n

Проверку ведут по одному jobId. В журнале достаточно событий received, result_saved, retry_scheduled, quarantined и ack_sent. Рядом пишут attempt, признак redelivered и безопасную причину. Полный payload, токены и пользовательские документы в журнал не кладут.

\n

Retry нужен не для любой ошибки

\n

Повтор оправдан, если новое время может изменить исход: зависимость временно недоступна, сработал сетевой timeout до ответа или ожидаемая запись ещё не появилась. Но timeout после отправки запроса не доказывает, что внешний эффект не состоялся. Такой вызов повторяют только с ключом идемпотентности или после проверки статуса операции.

\n

Невалидный JSON, неизвестная версия события и отсутствующее обязательное поле повтором не исправятся. Бесконечный nack(requeue=true) создаёт горячий redelivery loop: RabbitMQ может быстро вернуть сообщение в очередь, и consumer будет снова брать ту же запись. Он занимает worker и прячет полезные сообщения за одной постоянной ошибкой.

\n

Для retry задают максимальное число попыток, причину последнего перехода и время следующего допуска. Задержка в этой модели создаётся retry-очередью с TTL, планировщиком или другим механизмом выбранного клиента: dispatcher публикует новую delivery после notBefore, а текущую delivery отклоняют с requeue=false. Это отличается от немедленного requeue. Попытка хранится в базе или в retry-сообщении, потому что AMQP delivery сама по себе не является счётчиком попыток.

\n

Карантин для poison message

\n

Poison message — сообщение, которое текущий consumer не может обработать автоматически. Worker сначала сохраняет причину и состояние quarantined, затем отклоняет delivery без requeue. При настроенном dead-letter exchange broker переопубликует сообщение в отдельный exchange. Без такой конфигурации оно может быть отброшено, поэтому карантин должен быть проверяемой частью инфраструктуры, а не только словом в коде. Состояние в базе и сообщение в DLX — разные следы: нужно проверить оба.

\n

Карантин не означает успех. Он означает, что автоматический путь остановился с понятной причиной. Владелец может исправить payload и переиздать задачу, обновить consumer или отменить операцию. Автоматически читать карантин обратно в основную очередь без исправления причины нельзя: loop вернётся.

\n

Порядок действий

\n
  1. Определить стабильный jobId, область уникальности idempotency key и записывать их в задачу, событие, результат и журнал.
  2. Разделить состояния queued, running, retry_wait, succeeded и quarantined.
  3. Проверить согласованное создание outbox и задачи, а также публикацию после подтверждения broker-клиента.
  4. Поставить атомарный claim или lease и проверку terminal state до необратимого эффекта.
  5. Сохранить результат и resultKey до ack конкретного delivery на том же канале.
  6. Разделить временные, неопределённые и постоянные ошибки; для каждой задать проверку и лимит.
  7. Для временной ошибки записать retry-intent, отклонить текущую delivery с requeue=false и проверить задержанную публикацию.
  8. Для постоянной ошибки записать quarantined, отправить отказ без requeue и проверить dead-letter маршрут.
  9. Прогнать повторную доставку после сохранённого результата и убедиться, что второй эффект не создаётся.
\n

Ограничения

\n

Эта схема не делает систему ровно-однократной. Два worker могут одновременно выполнить эффект, если claim, lease или уникальное ограничение реализованы неверно. Внешний сервис может принять запрос и не вернуть ответ. Broker может иметь другую семантику подтверждений. Поэтому порядок нужно сверить с версией клиента, типом очереди и реальной политикой dead-lettering.

\n

Псевдокод выше не открывает соединение с RabbitMQ, не измеряет throughput и не является production-тестом. Он ограничен учебной иллюстрацией переходов. Интеграционная проверка должна использовать выбранный broker, несколько worker, падение после сохранения результата и до ack, повтор после результата, отдельный retry с задержкой и невалидный payload после лимита.

\n

Критерий готовности

\n

Решение готово, когда для одного заранее известного jobId журнал показывает устойчивый результат до первого ack, повторную delivery с признаком redelivered и отсутствие второго эффекта. Для временной ошибки видны сохранённый retry-intent, ограниченные попытки и следующий допуск. Для невалидного payload видны причина, состояние quarantined, отказ с requeue=false и подтверждённый маршрут DLX либо явно зафиксированное отбрасывание. Эти свойства должны воспроизводиться на интеграционном стенде выбранного broker, а не только в unit-тесте.

\n

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

\n" }