Симптом выглядит как «экспорт иногда исчезает». Пользователь нажал кнопку, 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 уже начал работу.
| Слой | Что хранит | Что считается доказательством | Чего не обещает |
|---|---|---|---|
| HTTP | jobId и ответ 202 | в транзакции появилась запись задачи | что обработка уже завершена |
| База приложения | state, attempt, resultKey, lastError | задача читается по jobId | что сообщение дошло до broker |
| Outbox | событие публикации и признак отправки | есть элемент, который можно дочитать после рестарта | атомарный commit с внешним broker |
| Очередь | delivery выбранному consumer | ручной ack или nack по delivery tag | идемпотентность бизнес-операции |
| Worker | результат и переход состояния | resultKey сохранён до ack | безопасность внешнего побочного эффекта без ключа |
Такой список убирает опасное сокращение «задача в очереди». У нас есть как минимум запись намерения, запись на публикацию, delivery и бизнес-результат. Они могут находиться в разных состояниях одновременно. Дежурный не должен гадать по отсутствию файла: он открывает строку задачи, смотрит state, attempt и lastError, после чего знает, на какой границе продолжать проверку.
Разрыв между записью и публикацией надо назвать
Наивная последовательность «сначала записали заявку, потом отправили сообщение» ломается при падении между двумя строками. В базе уже есть 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 на один и тот же отказ, а полезные задачи окажутся позади него.
| Наблюдение | Состояние задачи | Действие с delivery | Что остаётся для разбора |
|---|---|---|---|
| Сохранён CSV и resultKey | succeeded | ack | jobId, ключ результата, попытка |
| Короткий timeout зависимости, попытка 1–2 | retry_wait | вернуть в задержанный retry-маршрут | код ошибки и следующая попытка |
| Невалидный payload или попытка 3 | quarantined | nack/reject без requeue; DLX при настроенном маршруте | payload-версия, причина, jobId, delivery metadata |
| Повтор delivery после сбоя до ack | succeeded уже есть | ack без нового экспорта | признак redelivered и прежний resultKey |
Poison message здесь не мистическая категория broker. Это сообщение, которое снова и снова не может пройти известный consumer-контракт. Его путь должен заканчиваться в отдельной карантинной очереди или другой управляемой поверхности, а не возвращаться в тот же hot loop. Карантин не означает «удалить и забыть»: для записи должны остаться jobId, версия payload, причина, число попыток и понятный владелец решения — исправить данные, исправить worker или осознанно отменить работу.
Маршрут проверки до первого реального запуска
- Записать для операции постоянный
jobId, разрешённые состояния и условие, после которого пользователь может увидеть результат. - В одной транзакции приложения создать job и outbox; после имитации рестарта проверить, что непросланный outbox всё ещё читается.
- На учебной фикстуре прогнать временную ошибку: первая попытка переходит в
retry_wait, вторая сохраняет результат, и только затем появляетсяack_sent. - Отдельно прогнать невалидный payload или лимит попыток: запись становится
quarantined, а маршруту не разрешён немедленный requeue. - Повторить delivery для уже
succeededjob и убедиться, что результат не создаётся второй раз, а worker подтверждает новое delivery. - До подключения 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