diff --git a/editorial/agent-rewrites/245.json b/editorial/agent-rewrites/245.json index 9ade09c..bbdc652 100644 --- a/editorial/agent-rewrites/245.json +++ b/editorial/agent-rewrites/245.json @@ -1,7 +1,7 @@ { "index": 245, "slug": "editorial-2021-03-mechanism-queues", - "title": "Очереди задач: почему повторная доставка не должна повторять эффект", - "excerpt": "Сообщение может прийти повторно после уже выполненной операции. Разбираем границу между delivery и доменным эффектом, идемпотентный ключ, порядок retry и безопасный terminal route.", - "contentHtml": "
Симптом заметен не в очереди, а в результате: клиент получил два письма, счёт создался дважды или один заказ перешёл в неверный статус. В логах при этом видны два запуска одного обработчика. Команда часто обвиняет broker и пытается отключить повторную доставку. Это опасное решение. Вместе с повтором можно потерять задачу, если первый consumer успел выполнить часть работы, но не успел подтвердить delivery. Цена ошибки зависит от эффекта: лишнее уведомление можно отменить, второе списание — уже нет.
\nТезис статьи простой: подтверждение относится к текущей доставке сообщения, а идемпотентность относится к доменному эффекту. Эти границы нужно проектировать отдельно. Consumer должен переживать повтор одного намерения, хранить ключ эффекта и принимать решение о retry после проверки причины. Выбранная очередь помогает доставить работу, но не делает внешнюю операцию атомарной и не обещает exactly-once для всей системы.
\nУ одной задачи есть как минимум три разных идентификатора и состояния. messageId обозначает конкретное сообщение в транспорте. effectKey обозначает доменный эффект, например отправку напоминания по счёту. sequenceKey обозначает сущность, для которой важен порядок: счёт, заказ или профиль.
Broker отвечает за доставку сообщения и его подтверждение. Consumer отвечает за выполнение операции. Доменное хранилище отвечает за факт эффекта. Если процесс упал после записи в хранилище, но до ack, broker имеет право доставить сообщение ещё раз. Если код связывает повтор с новым эффектом, он превращает штатное восстановление в дубль.
\ndelivery: messageId = msg-81, attempt = 2\neffect: effectKey = invoice-417:reminder, state = recorded\norder: sequenceKey = invoice-417, next = 8\nack: acknowledge current delivery after effect decision\nТакая модель не говорит, что повтор всегда произойдёт. Она говорит, что код не должен ломаться, если подтверждение потерялось, соединение закрылось или consumer завершился в неудобный момент. При ручном подтверждении неподтверждённая доставка обычно возвращается в работу после закрытия канала или соединения. Поэтому запись эффекта должна предшествовать ack, а проверка повторного эффекта — предшествовать новой записи.
\nРассмотрим учебный пример без подключения к broker. Сообщение просит отправить одно напоминание по счёту. Первая попытка вызывает временную ошибку и уходит на retry. Вторая попытка записывает эффект в ledger. Сразу после записи процесс теряет соединение. Consumer не знает, дошёл ли ack. Broker считает доставку неподтверждённой и запускает её снова.
\nНа третьем запуске новый код сначала ищет effectKey. Ledger возвращает существующую запись. Consumer не отправляет новое напоминание и завершает текущую доставку как обработанную. В логах остаётся факт duplicate, но в домене появляется одна запись. Это ожидаемый результат recovery, а не доказательство exactly-once.
async function handle(message, ledger) {\n const key = message.effectKey;\n const existing = await ledger.find(key);\n\n if (existing) {\n await ledger.recordDelivery(message.messageId, 'duplicate-effect-suppressed');\n return { ack: true, effect: 'not-repeated' };\n }\n\n await ledger.recordEffect({\n effectKey: key,\n messageId: message.messageId,\n type: message.type,\n });\n\n return { ack: true, effect: 'recorded' };\n}\nКод учебный. Он показывает порядок решений, но не заменяет транзакцию, уникальный индекс или API внешнего сервиса. В рабочей системе find и recordEffect должны защищать одну границу состояния. Иначе два параллельных consumer могут одновременно не найти ключ и оба создать эффект. Для базы это обычно означает уникальное ограничение по effectKey и обработку конфликта как duplicate. Для внешнего HTTP-вызова нужен поддержанный внешней системой idempotency key или ручной контроль. Локальный Map не может отменить уже отправленное письмо.
| Симптом | Причина | Проверка | Действие |
|---|---|---|---|
| Один эффект записан дважды | Нет уникального effectKey или проверка не атомарна | Сравнить ключи и найти две записи в одном интервале | Добавить уникальное ограничение и трактовать конфликт как duplicate |
| Задача пропала после падения worker | Ack отправили до записи эффекта или включили auto-ack | Сопоставить время ack, запись эффекта и завершение процесса | Подтверждать после успешной границы обработки; для неизвестного исхода включить recovery |
| Очередь быстро растёт | Retry повторяет постоянную ошибку или consumer не успевает | Разделить transient и permanent причины, посмотреть attempts и latency | Задать лимит попыток, backoff и terminal route |
| Статус вернулся назад | Параллельные сообщения нарушили порядок для одной сущности | Сгруппировать события по sequenceKey и сравнить номера | Проверять следующий номер; gap отправлять на разбор, а не угадывать |
| Оператор повторил опасную задачу вслепую | Terminal запись не содержит причины и ключа эффекта | Проверить содержимое ручного маршрута | Сохранять messageId, effectKey, attempts, reason и requiredCheck |
Retry подходит для ограниченного класса отказов: временно недоступна зависимость, закончился connection pool или сработал rate limit. Он не исправляет неправильный формат сообщения, отсутствующий обязательный атрибут или нарушение бизнес-правила. Такой вход будет падать снова. Без лимита consumer создаст requeue loop, нагрузит broker и отложит диагностику.
\nПолитика должна различать причину, число попыток и следующий исход. Backoff снижает плотность повторов, но не сообщает, когда ошибка стала постоянной. После лимита попыток сообщение нужно перевести в явный terminal route: dead-letter queue, quarantine или ручной разбор. Название зависит от продукта. Контракт должен оставаться одинаковым: задача перестала исполняться автоматически, а причина и контекст сохранены.
\nconst policy = {\n transient: { delaysMs: [1000, 4000, 16000], terminal: 'manual-review' },\n permanent: { delaysMs: [], terminal: 'quarantine' },\n};\n\nfunction nextAction(error, attempt) {\n const rule = error.kind === 'transient' ? policy.transient : policy.permanent;\n if (attempt < rule.delaysMs.length) return { type: 'retry', delayMs: rule.delaysMs[attempt] };\n return { type: 'terminal', route: rule.terminal };\n}\nЗначения в примере учебные. Их нельзя переносить в production без проверки SLA зависимости, лимитов broker и допустимого времени ожидания. Отдельно измеряйте число повторов, возраст самой старой задачи и долю terminal исходов. Среднее время обработки может выглядеть нормальным, пока небольшой поток poison messages держит ресурсы и скрывает реальную причину.
\nОчередь не обязана сохранять общий порядок всех задач. Обычно нужен порядок только внутри одного ключа. Для invoice-417 событие с номером 8 можно применить после номера 7. Событие с номером 10 нельзя молча применить раньше 9, если доменная модель не допускает пропуск. Для разных счетов искусственная последовательность только уменьшит параллелизм.
Порядок должен иметь владельца и проверяемое правило. Partition, routing key или один worker могут помочь доставке, но не заменяют проверку состояния. После перезапуска consumer должен снова понять, какой номер уже принят. Если предыдущего события нет, выберите один из исходов: подождать ограниченное время, запросить восстановление или отправить gap в manual route. Бесконечный retry здесь маскирует потерю данных.
\nmessageId, effectKey и при необходимости sequenceKey. Зафиксируйте, какие сообщения законно создают разные эффекты.Эта схема не делает распределённую систему атомарной. Если запись в локальной базе и вызов внешнего API идут в разных системах, между ними остаётся окно неопределённости. В нём возможны внешний успех без локальной записи, локальная запись без внешнего успеха и повтор после сетевого таймаута. Решение выбирают по доменному риску: outbox, API с идемпотентным ключом, сверка состояния или ручная операция. Ни один вариант не следует объявлять универсальным без проверки конкретных границ.
\nПорядок тоже ограничен областью ключа. Один partition или один consumer не создаёт общий порядок между независимыми сущностями. Флаг redelivered не доказывает, что сообщение ранее полностью обработали, и отсутствие такого флага не доказывает обратное. Наблюдайте историю delivery, но принимайте решение по состоянию эффекта.
\nСчитайте контракт готовым, когда контролируемый тест проходит один и тот же сценарий: consumer записывает эффект, теряет знание об ack, получает повторное сообщение и оставляет ровно один доменный эффект; лог и ledger связывают оба запуска с одним effectKey; постоянная ошибка после лимита попадает в terminal route с причиной; gap не меняет состояние раньше времени. Тест должен выполняться на выбранном broker, хранилище и клиенте проекта. Учебный пример выше проверяет только порядок решений и не заменяет эту интеграционную проверку.
Симптом выглядит как спор с логами: обработчик записал эффект, но затем то же сообщение пришло снова. Если второй запуск создаёт ещё одно письмо, списание или запись, команда начинает искать ошибку в очереди. Цена такой реакции выше дубля. Можно настроить повтор иначе и всё равно оставить ту же дыру между эффектом и подтверждением. Пока не названа граница, где эффект считается записанным, любой термин о delivery скрывает главный вопрос: что делать с повторной работой над тем же намерением.
\nДля марта 2021 года полезно говорить скромнее. Подтверждение доставки — это решение по конкретной доставке сообщения. Идемпотентность — свойство операции с определённым ключом эффекта. Порядок — инвариант доменного ключа. Эти вещи могут взаимодействовать, но не становятся одним свойством после выбора broker. Ниже — детерминированная state machine в памяти. Она не соединяется с AMQP или Kafka и не заявляет, что даёт exactly-once. Её задача — сделать видимыми точки, где появляется duplicate и где он должен быть остановлен.
\nСлова at-most-once, at-least-once и exactly-once часто попадают в решение раньше контракта. Для локального дизайна полезнее разложить их на наблюдаемые обязательства. Мы можем договориться, что consumer допускает повтор одной логической задачи; что эффект с одинаковым effectKey записывается один раз; что порядок проверяется только по sequenceKey; что неизвестная ошибка не повторяется бесконечно. Ни одно из этих предложений не превращает абстрактную модель в характеристику сети, диска или выбранного продукта.
| Термин в обсуждении | Что фиксируем в учебном контракте | Что остаётся за границей | Проверяемый признак |
|---|---|---|---|
| delivery | один запуск consumer над message id | дошло ли сообщение по сети и как хранит его broker | в fixture видны delivery 1, 2 и 3 |
| duplicate | повторный запуск того же id после неопределённого результата | почему именно появился повтор в конкретном transport | ledger возвращает duplicate-effect-suppressed |
| effect | доменная запись, привязанная к effectKey | атомарность между базой и внешней системой | в ledger остаётся одна строка ключа |
| order | проверка последовательности для одного sequenceKey | общий порядок всех задач | gap ведёт к manual review, а не к догадке |
| terminal route | автоматический маршрут остановлен с контекстом | решение оператора и последующая интеграция | есть reason, attempts и requiredCheck |
Такой словарь снимает ложный выбор между «настроить гарантию» и «ничего не делать». У команды появляется ряд маленьких вопросов. Когда можно подтвердить обработку? Где лежит ключ эффекта? Что происходит, если внешний вызов завершился, а запись о нём нет? Для какой сущности порядок обязателен? Какой случай не имеет права автоматически возвращаться в работу? Ответы могут оказаться разными даже внутри одного сервиса. Это нормально: граница договора определяется риском операции, а не названием очереди.
\nПредставим учебный второй запуск. Первая попытка получила временное условие и была отложена на 1000 мс. Вторая дошла до записи эффекта и сразу после этого потеряла знание о результате подтверждения. Третья доставка выглядит как duplicate. Если обработчик не смотрит в ledger, он повторит эффект. Если он смотрит в ledger до эффекта, он видит существующий ключ и может завершить текущую доставку без нового доменного действия. Эта последовательность не доказывает, что повтор обязательно случится; она показывает, почему код обязан быть готов к нему.
\nfunction recordEffectOnce(ledger, message) {\n if (ledger.has(message.effectKey)) {\n return { state: 'duplicate-effect-suppressed' };\n }\n ledger.set(message.effectKey, { messageId: message.id });\n return { state: 'effect-recorded' };\n}\n\n// Подтверждение доставки принимается только после решения по ledger.\nВажна последовательность, а не название функции. Сначала проверить effectKey. Если ключ найден, не создавать второй эффект. Если ключа нет, попытаться записать эффект и ключ в одной подходящей для проекта границе. После этого принять решение о подтверждении текущей доставки. Чем дальше друг от друга эти действия, тем больше сценариев неопределённости. Внешний HTTP-вызов особенно важен: Map из фикстуры не способна откатить письмо или платёж. Там нужен отдельный контракт идемпотентного ключа на стороне внешней границы либо ручный маршрут.
function handleTrainingDelivery(ledger, message) {\n const effect = recordEffectOnce(ledger, message);\n if (effect.state === 'duplicate-effect-suppressed') {\n return { delivery: 'ack', effect: 'not-repeated' };\n }\n return { delivery: 'ack', effect: 'recorded-once-in-training-ledger' };\n}\n\n// В production эта функция должна получить реальную границу хранения.\nКогда ledger подавил повтор, обработка ещё не закончила объяснение. Нужно сохранить, что повтор был и на каком участке он обнаружен. Иначе через месяц останется только одна строка эффекта, но пропадёт сигнал, что граница подтверждения или восстановление consumer требуют отдельной проверки. В учебной fixture третий delivery имеет duplicate: true и результат ack-after-ledger-check. Это не показатель реального флага выбранного протокола, а документированный исход модели.
Полезный контрпример: не хранить один общий список message id без связи с эффектом. Если один logical id законно создаёт несколько различных эффектов, глобальный список помешает работе. Если два разных message id представляют одно и то же намерение, список id не остановит duplicate. Поэтому ключ выбирают у доменного эффекта: invoice-417:reminder в учебном примере означает именно одно напоминание для конкретного счёта. Это решение не универсально; название и состав ключа должен подтвердить владелец домена.
const fixture = runQueueFixture();\nif (!Object.values(fixture.assertions).every(Boolean)) {\n throw new Error('training queue contract failed');\n}\n\nconsole.log(fixture.recovery.map((item) => item.outcome));\n// ['controlled-retry', 'effect-written-receipt-unknown', 'ack-after-ledger-check']\nФикстура проверяет десять инвариантов: один message id в обеих учебных ветвях, контролируемую задержку, отсутствие эффекта на retry, единственную запись ledger, suppress duplicate, terminal решение по duplicate, пустой ledger в poison path, manual route, границу решения оператора и разделение двух ветвей. Она не измеряет retry клиента и не показывает протокольный acknowledgement. Её ценность в том, что при редактуре или доработке нельзя тихо поменять правило на «повтор создаёт новый эффект».
\nВопрос порядка часто появляется поздно: сначала consumer обработал несколько задач параллельно, затем доменная модель требует, чтобы статус не вернулся назад. Здесь недостаточно сказать «сделаем один worker». Один worker замедлит всё, но не объяснит, что происходит после перезапуска или между разными ключами. Нужен sequenceKey, правило следующего номера и исход для gap. Для несвязанных задач правило может отсутствовать; искусственный порядок там превращает обработку в очередь ожидания без пользы.
const lane = new Map();\n\nfunction acceptInSequence(message) {\n const previous = lane.get(message.sequenceKey) || 0;\n if (message.sequence !== previous + 1) {\n return { state: 'manual-review', reason: 'sequence-gap' };\n }\n lane.set(message.sequenceKey, message.sequence);\n return { state: 'ready' };\n}\n\n// Это локальный контракт одного ключа, не обещание глобального порядка.\nЭтот пример не реализует partition или блокировку. Он показывает форму инварианта: для invoice-417 можно принять номер только после предыдущего. Если номер пропущен, consumer не придумывает порядок и не раздувает retry. Он формирует terminal record с причиной sequence-gap. Дальше владелец домена смотрит на происхождение входа: задача пришла раньше, потеряна запись о предыдущем шаге или последовательность вообще неправильно определена. Это уже другое расследование, не алгоритм повторной доставки.
Задержка перед повтором нужна, чтобы не превращать временный сбой в плотный цикл. Но backoff не делает ошибку временной и не восстанавливает порядок. В учебной policy две задержки заданы явно: 1000 и 4000 мс. Они не вычисляются из случайного времени и поэтому fixture повторяема. В реальном проекте значения выбирают по договору зависимости, допустимому ожиданию и наблюдаемой нагрузке. До такого выбора надо разделить хотя бы две причины: контролируемое временное условие и вход, который consumer не умеет обработать.
\nconst retryPolicy = { maxAttempts: 2, delaysMs: [1000] };\n\nfunction nextTrainingDecision(attempt, failureKind) {\n if (failureKind === 'temporary' && attempt < retryPolicy.maxAttempts) {\n return { state: 'retry', afterMs: retryPolicy.delaysMs[attempt - 1] };\n }\n return { state: 'manual-review', reason: failureKind };\n}\n\nnextTrainingDecision(1, 'temporary');\n// { state: 'retry', afterMs: 1000 }\nЕсли retry уже исчерпан, терминальный путь должен быть видимым, а не состоять из удаления сообщения. Ручная запись несёт id, effectKey, причину, попытки и то, какую проверку ожидают от оператора. Это не бюрократия. Без effectKey оператор не знает, может ли replay создать второе действие. Без причины неясно, исправлять ли вход или зависимость. Без requiredCheck любой повтор становится случайным запуском того же consumer.
Спецификация AMQP 0-9 отделяет delivery и acknowledgement, а также содержит признак redelivered. Kafka 2.7 в своей versioned documentation отдельно обсуждает семантику доставки и последствия retry. Эти факты полезны, потому что не дают свести обработку к слову «очередь». Но в статье нет вывода о реальном конфиге RabbitMQ или Kafka: учебные имена, Map и события fixture не соответствуют API какого-либо продукта. Переносить нужно вопросы к контракту, а не код из примера.
\nЗдесь не запускались broker, HTTP, база, внешнее API, browser, CI или production build. Нет измерений throughput, времени восстановления или потери данных. Следующий шаг после чтения — показать выбранную транзакционную границу и ключ эффекта на маленькой интеграции. Если эту границу невозможно обеспечить, надо сократить автоматическое действие и направить сомнительные случаи в ручной маршрут, а не назвать задачу solved из-за одного успешного запуска.
\n