DarkRiDDeR15 мин

Очереди задач: какой контракт написать до первого consumer

НадёжностьBackendПрактика

Симптом появляется после первого успешного consumer: напоминание ушло дважды, соседнее изменение пришло раньше предыдущего, а запись с неподходящей схемой снова возвращается в обработку. Цена не в самом факте очереди. Цена в том, что команда не может сказать, какой эффект уже сделан, кому принадлежит порядок и где остановить сообщение, которое нельзя обработать. Тогда повтор становится случайной попыткой, а ручной разбор начинается без исходного контекста.

В марте 2021 года я бы не начинал с обещания, что один broker решит доставку. Сначала нужен контракт на одну логическую задачу. Он определяет, что наблюдаем как сообщение, что считаем внешним эффектом, когда допускаем повтор и какой исход переводит задачу в ручной маршрут. Ниже все значения учебные и живут в памяти Node. Это не RabbitMQ, Kafka, SQS, реальный consumer или замер нагрузки. Модель нужна, чтобы проверить границы решения до подключения инфраструктуры.

Очередь отделяет выполнение, но не отменяет договор

Очередь полезна, когда создатель задачи не должен ждать работу обработчика. Но она добавляет границу между тем, кто сформировал вход, и тем, кто создаёт эффект. По эту сторону границы есть запись о намерении. По другую — отправленное письмо, изменённый статус, созданный файл или вызов соседнего API. Если эти два состояния не связаны явным ключом, повторная доставка легко превращается в повторный эффект. Поэтому первый вопрос не «какой consumer написать», а «какой результат он имеет право создавать и как мы увидим, что результат уже был».

Порядок тоже нельзя приписывать очереди одним словом. Для части задач порядок не важен: две независимые очистки можно выполнить в любом порядке. Для другой части есть один объект, например счёт или заказ, и изменение с номером 12 нельзя применять до номера 11. Это локальный инвариант одного ключа, а не глобальная очередь для всей системы. У такого инварианта должен быть владелец: функция, которая хранит последнюю принятую последовательность, или явный маршрут на ручную проверку при пропуске.

Минимальный контракт задачи: поле должно отвечать на конкретный вопрос
Поле или правилоВладелецИнвариантПроверка перед действием
idсоздатель задачиодна логическая задача имеет один наблюдаемый idid не пустой и сохраняется во всех учебных ветвях
effectKeyобработчик и ledgerодин ключ даёт не более одного записанного эффектасначала lookup в ledger, затем запись или suppress duplicate
sequenceKeyдоменный владелец порядкапорядок проверяется только внутри одного ключаследующий номер следует за предыдущим или задача уходит в manual route
retry policyконтракт обработкиповтор ограничен числом попыток и задержкойвременный класс ошибки соответствует известному условию
manual routeоператор и владелец доменатерминальный случай не исчезает и не зацикливаетсяесть причина, число попыток и допустимые следующие действия

Эта таблица не требует отдельной платформы. Её можно положить рядом с функцией обработчика и с тестом контракта. Главное — не смешивать роли. id позволяет найти историю. effectKey защищает конкретный эффект. sequenceKey ограничивает порядок. Попытка использовать один случайный идентификатор для всех трёх задач обычно скрывает важные вопросы: что делать при новой версии входа, можно ли повторить операцию после исправления и относится ли порядок к одному объекту или к очереди целиком.

Конверт сообщения описывает границу до payload

Payload отвечает за данные доменной операции. Конверт отвечает за то, как эту операцию вести через границу обработки. В учебном конверте есть id, kind, effectKey, sequenceKey и минимальная версия схемы. Этого мало для любого приложения, но достаточно, чтобы прочитать запись о повторе: это тот же логический запрос или новый; он должен создавать тот же эффект или другой; порядок нужен именно для этого счёта или его здесь вообще нет.

// Минимальный контракт сообщения. Значения учебные.
const message = {
  id: 'msg-order-417',
  kind: 'invoice.reminder',
  effectKey: 'invoice-417:reminder',
  sequenceKey: 'invoice-417',
  payload: { invoiceId: '417', schema: '2021-03' },
};

// id отвечает за наблюдение, effectKey — за внешний эффект.
// sequenceKey нужен только там, где порядок действительно является инвариантом.

Схема не должна быть декоративной строкой. Если consumer не знает, как трактовать payload, он не должен угадывать поле и продолжать работу с частично понятным объектом. В нашем учебном примере неизвестная версия даёт manual-review. Такой исход не говорит, что запись испорчена навсегда. Он говорит ровно одно: текущий обработчик не имеет безопасного правила для этого входа. Исправление может оказаться в producer, в миграции схемы или в новом consumer, но это отдельное решение с новой проверкой.

function validateTrainingMessage(message) {
  if (!message.id) throw new Error('message id is required');
  if (!message.effectKey) throw new Error('effect key is required');
  if (!message.payload || message.payload.schema !== '2021-03') {
    return { state: 'manual-review', reason: 'schema-not-supported' };
  }
  return { state: 'ready', sequenceKey: message.sequenceKey };
}
Вертикальная схема жизненного цикла учебного сообщения: producer создаёт конверт, consumer валидирует его, временная ошибка идёт в ограниченный retry, effect ledger подавляет duplicate, а терминальная ошибка сохраняет запись для manual review
Жизненный цикл учебной задачи: retry — один из исходов, а не бесконечная ветка по умолчанию.

Повтор исправляет только временное условие

Повтор полезен, когда причина уже названа временной и есть проверяемый предел. Например, в учебной модели зависимость не ответила на первый контролируемый вызов, поэтому задача ждёт 1000 мс и возвращается на вторую попытку. Эта задержка не доказывает, что в реальной системе зависимость восстановится, и не является рекомендацией для любой нагрузки. Она делает другое: ограничивает состояние модели. После заданного числа попыток вопрос меняется с «когда повторить» на «какой человек или код должен принять решение дальше».

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

const retryPolicy = { maxAttempts: 2, delaysMs: [1000] };

function nextTrainingDecision(attempt, failureKind) {
  if (failureKind === 'temporary' && attempt < retryPolicy.maxAttempts) {
    return { state: 'retry', afterMs: retryPolicy.delaysMs[attempt - 1] };
  }
  return { state: 'manual-review', reason: failureKind };
}

nextTrainingDecision(1, 'temporary');
// { state: 'retry', afterMs: 1000 }

В коде nextTrainingDecision принимает классификацию извне. Это намеренно. Функция не может по строке исключения честно узнать, временная ли причина в любой системе. Отдельная граница отвечает за классификацию и обязана вернуть понятный вид причины. Когда этого правила нет, лучше записать неуверенный случай в ручной маршрут, чем назвать каждую ошибку temporary и ждать, пока ограничение на попытки само скроет проблему.

Ledger принадлежит эффекту, а не слову retry

Есть неприятный промежуток: эффект уже записан, но consumer не знает, дошло ли его подтверждение. Следующая доставка может быть тем же логическим сообщением. Если код сначала отправит письмо или изменит запись, а затем попробует понять, был ли такой эффект раньше, duplicate уже случится. Поэтому проверка ключа эффекта должна идти перед созданием эффекта, а результат этой проверки должен позволять завершить повтор без второго действия.

function recordEffectOnce(ledger, message) {
  if (ledger.has(message.effectKey)) {
    return { state: 'duplicate-effect-suppressed' };
  }
  ledger.set(message.effectKey, { messageId: message.id });
  return { state: 'effect-recorded' };
}

// Подтверждение доставки принимается только после решения по ledger.

В реальном проекте ledger может быть строкой в той же транзакции, уникальным индексом или другим доменным механизмом. Этот выбор зависит от базы, внешнего API и границы транзакции. Учебная Map ничего такого не реализует. Она показывает единственный инвариант: второй вызов с тем же effectKey не добавляет новую запись. Нельзя называть это гарантией доставки. Это защита конкретного доменного эффекта при условии, что реальная запись ledger и сам эффект имеют согласованную границу.

Маршрут до первого запуска

  1. Назвать один тип задачи и один внешний эффект. Не объединять отправку письма, обновление баланса и перестройку индекса в один общий контракт.
  2. Добавить в конверт стабильный id, отдельный effectKey и sequenceKey только при реальной необходимости порядка.
  3. Сформулировать допустимые исходы consumer: успешно записан эффект, duplicate подавлен, controlled retry, manual review. Для каждого указать владельца.
  4. Записать retry policy с числом попыток, задержкой и одним классом временной причины. Все неясные случаи оставить вне автоматического повтора.
  5. Выбрать место ledger рядом с реальным эффектом и проверить, что duplicate не может создать вторую запись или второй вызов.
  6. Сохранить terminal record с id, effectKey, причиной, попытками и требуемым ручным решением. Не заменять её строкой в логе без владельца.
  7. Запустить контролируемую fixture, затем отдельно проверить выбранный broker, хранилище и внешний API в среде проекта.

Историческая рамка и границы примера

К марту 2021 года AMQP 0-9 уже содержал отдельный признак повторной доставки и подтверждение доставки. В обзоре RabbitMQ 3.8, выпущенном в ноябре 2019 года, уже описывался delivery limit для poison message. Документация Kafka 2.7 тоже разводит сценарии доставки и последствия retry. Эти источники помогают выбрать вопросы, но не дают один переносимый конфиг. У каждого продукта свои версии, подтверждения, топология и граница между обработкой и внешним эффектом.

Фикстура этой статьи не открывает сеть, не публикует сообщение, не запускает broker и не создаёт реальный delivery guarantee. В ней один неизменяемый message id показан в двух взаимоисключающих учебных ветвях: recovery объясняет duplicate после неизвестного статуса подтверждения, poison объясняет terminal manual route. Это не одна фактическая история доставки. Следующий проверяемый шаг — взять свой тип задачи, написать такой же контракт и прогнать его на тестовой интеграции с версиями конкретного стека.

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

  • AMQP 0-9: спецификация basic.deliver, basic.ack и redelivered — первичная спецификация 2008 года: у delivery есть признак повторной доставки, а подтверждение относится к доставленному сообщению; учебная модель ниже не реализует протокол
  • RabbitMQ 3.8 Release Overview — 11 ноября 2019 года — версия и обзор были доступны к марту 2021 года; источник упоминает delivery limit для poison message, но не является описанием данного учебного маршрута
  • Apache Kafka 2.7.X: versioned documentation — версионная документация линии 2.7, доступной в марте 2021 года; терминология доставки приводится только для разграничения обязательств, не как claim об этой fixture
  • Apache Kafka 2.7.0 KafkaProducer API — официальный API линии 2.7 предупреждает, что retries могут открыть путь к duplicate; это не настройка и не запуск Kafka в статье