DarkRiDDeR14 мин

Фоновые задачи: сначала фиксируем намерение, потом запускаем worker

BackendОчередиПрактика

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

В мае 2020 я бы не начинал с «фоновой платформы». Достаточна узкая связка delivery и backend: HTTP принимает намерение, приложение сохраняет состояние задачи, отдельный dispatcher публикует сообщение, worker делает работу, а очередь получает подтверждение только после сохранённого результата. Это учебная схема для одной операции report.export. Она не доказывает работу конкретного RabbitMQ-кластера и не обещает exactly-once; её задача — сделать каждую точку потери или повтора видимой.

У задачи есть свой владелец состояния

Сообщение в очереди не должно быть единственным местом, где живёт смысл работы. Брокер знает о delivery, но не обязан знать, создан ли CSV, обновлена ли строка в базе и что увидит пользователь. Поэтому у операции есть запись приложения с постоянным jobId, входными параметрами, числом попыток, текущим состоянием, ключом результата и короткой причиной последней неудачи. Один и тот же jobId проходит HTTP, outbox, message, worker и журнал. Это не трассировка всей системы, а минимальная нить для одного расследования.

Нормальный ответ HTTP — не «готово», а «принято»: { accepted: true, jobId }. Клиент затем спрашивает статус именно этой записи или получает уведомление привычным для продукта способом. Если запись не создана, нет принятой задачи. Если она создана, но сообщение ещё не опубликовано, это отдельное наблюдаемое состояние, а не повод сказать пользователю, что worker уже начал работу.

Контракт одной фоновой задачи: состояние не прячется внутри очереди
СлойЧто хранитЧто считается доказательствомЧего не обещает
HTTPjobId и ответ 202в транзакции появилась запись задачичто обработка уже завершена
База приложенияstate, attempt, resultKey, lastErrorзадача читается по jobIdчто сообщение дошло до broker
Outboxсобытие публикации и признак отправкиесть элемент, который можно дочитать после рестартаатомарный commit с внешним broker
Очередьdelivery выбранному consumerручной ack или nack по delivery tagидемпотентность бизнес-операции
Workerрезультат и переход состоянияresultKey сохранён до ackбезопасность внешнего побочного эффекта без ключа

Такой список убирает опасное сокращение «задача в очереди». У нас есть как минимум запись намерения, запись на публикацию, delivery и бизнес-результат. Они могут находиться в разных состояниях одновременно. Дежурный не должен гадать по отсутствию файла: он открывает строку задачи, смотрит state, attempt и lastError, после чего знает, на какой границе продолжать проверку.

Вертикальная схема жизненного цикла фоновой задачи: HTTP сохраняет задачу и outbox, dispatcher публикует сообщение, worker сохраняет результат, затем отправляет ack; ошибочная задача уходит в повтор или карантин
Жизненный цикл показывает четыре разных факта: намерение записано, сообщение опубликовано, результат сохранён, delivery подтверждён. Ack стоит последним, потому что не заменяет бизнес-результат.

Разрыв между записью и публикацией надо назвать

Наивная последовательность «сначала записали заявку, потом отправили сообщение» ломается при падении между двумя строками. В базе уже есть queued, а delivery нет. Обратный порядок не лучше: worker может получить сообщение и не найти ещё незафиксированную запись. Распределённую транзакцию между базой и broker я здесь не предлагаю: для небольшого сервиса она быстро становится дороже самой задачи. Вместо этого полезнее в одной транзакции приложения сохранить и задачу, и outbox-запись.

После commit отдельный короткий dispatcher ищет outbox без publishedAt, передаёт минимальное сообщение { jobId, type, version } выбранному клиенту broker и отмечает публикацию только по контракту этого клиента. Если процесс умер до такой отметки, dispatcher попробует снова. Отсюда следует неприятный, но здоровый вывод: consumer обязан выдержать duplicate. Повтор публикации — не дефект, который можно «выключить» флагом; это цена за возможность восстановиться после неясного обрыва.

// Учебный псевдокод: одна транзакция приложения, не клиент RabbitMQ.
async function requestExport(input, db) {
  const jobId = makeStableId(input.accountId, input.period);

  await db.transaction(async (tx) => {
    await tx.insertJob({
      id: jobId, state: 'queued', attempt: 0, resultKey: null,
    });
    await tx.insertOutbox({
      type: 'report.export.requested', jobId: jobId, publishedAt: null,
    });
  });

  return { accepted: true, jobId: jobId };
}

// Отдельный dispatcher читает неопубликованный outbox.
// Он отмечает publishedAt только после подтверждения выбранного broker client.

В примере нет SQL-схемы реального проекта и нет вызова библиотеки RabbitMQ. Важно другое: jobId создаётся до публикации, а outbox живёт рядом с бизнес-записью. В результате можно отдельно проверить, почему dispatcher не двигает запись: отсутствует ли соединение, неверен ли маршрут, не получено ли ожидаемое подтверждение или просто нет самого outbox-события. Это намного полезнее, чем повторно запускать весь экспорт по кнопке.

Ack подтверждает delivery, а не желание верить в успех

В AMQP delivery tag относится к конкретной доставке на конкретном канале, а не к вечному идентификатору задачи. Ручной ack говорит broker, что consumer принял ответственность за это delivery; после него broker может удалить сообщение. Поэтому ack до записи результата опасен: процесс может упасть после подтверждения, а нужная работа исчезнет из очереди. Ack после результата допускает обратное окно: результат есть, а ack не успел уйти. Тогда сообщение будет доставлено повторно. Это ожидаемый сценарий, не исключение.

Проверка порядка должна быть предельно приземлённой. До ack в хранилище уже есть state = succeeded и устойчивый resultKey. При повторной доставке worker читает эту запись, не запускает экспорт снова и подтверждает только повторное delivery. Если результат создаётся во внешнем сервисе, например объектном хранилище или email-шлюзе, одного флага succeeded недостаточно: нужен стабильный ключ объекта или ключ идемпотентности на стороне такого вызова. В этой статье мы ограничиваемся одной задачей и явно не выдаём этот принцип за универсальную гарантию всех сторонних систем.

Повтор — это отдельное состояние, не бесконечный requeue

У временной ошибки есть диагностируемая причина: короткий сбой зависимости, временный лимит или неготовый вход. Для неё можно сохранить retry_wait, номер попытки и код причины, а затем вернуть работу через заданную задержку и ограниченный маршрут. У невалидного входа причина другая: новый delivery не исправит неизвестный тип отчёта или отсутствующий обязательный параметр. Если каждое такое сообщение немедленно отправлять с requeue: true, worker будет тратить CPU на один и тот же отказ, а полезные задачи окажутся позади него.

Минимальная retry-политика для учебного экспорта
НаблюдениеСостояние задачиДействие с deliveryЧто остаётся для разбора
Сохранён CSV и resultKeysucceededackjobId, ключ результата, попытка
Короткий timeout зависимости, попытка 1–2retry_waitвернуть в задержанный retry-маршруткод ошибки и следующая попытка
Невалидный payload или попытка 3quarantinednack/reject без requeue; DLX при настроенном маршрутеpayload-версия, причина, jobId, delivery metadata
Повтор delivery после сбоя до acksucceeded уже естьack без нового экспортапризнак redelivered и прежний resultKey

Poison message здесь не мистическая категория broker. Это сообщение, которое снова и снова не может пройти известный consumer-контракт. Его путь должен заканчиваться в отдельной карантинной очереди или другой управляемой поверхности, а не возвращаться в тот же hot loop. Карантин не означает «удалить и забыть»: для записи должны остаться jobId, версия payload, причина, число попыток и понятный владелец решения — исправить данные, исправить worker или осознанно отменить работу.

Маршрут проверки до первого реального запуска

  1. Записать для операции постоянный jobId, разрешённые состояния и условие, после которого пользователь может увидеть результат.
  2. В одной транзакции приложения создать job и outbox; после имитации рестарта проверить, что непросланный outbox всё ещё читается.
  3. На учебной фикстуре прогнать временную ошибку: первая попытка переходит в retry_wait, вторая сохраняет результат, и только затем появляется ack_sent.
  4. Отдельно прогнать невалидный payload или лимит попыток: запись становится quarantined, а маршруту не разрешён немедленный requeue.
  5. Повторить delivery для уже succeeded job и убедиться, что результат не создаётся второй раз, а worker подтверждает новое delivery.
  6. До подключения broker выбрать конкретный клиентский API, проверить его publisher-confirm и DLX-настройки на стенде; этот текст не подменяет такой прогон.

Граница этой практики

Здесь нет настоящего очередного сервера, зарегистрированного consumer или production-лога. В модуле пакета есть только детерминированная in-memory фикстура состояний: она доказывает, что авторский переход «повтор → успех» отправляет ack после result и что третий неуспех переводит другую задачу в карантин. Она не доказывает поведение сети, задержку broker, порядок в нескольких worker и конфигурацию dead-letter exchange. Эти свойства должны быть проверены отдельным стендом с тем broker и клиентом, которые выбрал проект.

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

  • 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