Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Проектируйте потребителей сообщений так, чтобы многократная обработка одного и того же сообщения давала тот же эффект, что и его однократная обработка. Системы обмена сообщениями, гарантирующие доставку как минимум один раз, могут доставлять одно сообщение несколько раз. Без защиты от дубликатов повторная обработка сообщения может привести к появлению повторяющихся записей, к двойному списанию средств с клиента или к другим нежелательным последствиям.
Контекст и проблема
Распределенные приложения обычно обмениваются задачами через брокер сообщений вместо прямых синхронных вызовов. Большинство брокеров, включая Служебная шина Azure, Центры событий Azure, Apache Kafka и RabbitMQ, предоставляют по крайней мере один раз доставку. Это гарантирует, что сообщение достигает потребителя даже при возникновении сбоев, но также означает, что брокер может доставлять одно и то же сообщение более одного раза.
Дубликаты возникают из нескольких источников:
Повторные попытки продюсера
Производитель отправляет сообщение, не получает подтверждения из-за временных сбоев сети или времени ожидания и отправляет сообщение еще раз. Теперь у брокера есть две копии, хотя отправка с первой попытки прошла успешно.
Повторная доставка при отсутствии подтверждения
Потребитель получает и обрабатывает сообщение, но не подтверждает его получение, поскольку происходит сбой, истекает срок действия блокировки или теряется сообщение о подтверждении. Брокер предполагает, что сообщение не было обработано и снова доставляет его.
Сбои потребителей в процессе обработки
Потребитель завершает запись в базу данных, но аварийно завершает работу, прежде чем подтвердить получение сообщения. Когда другой экземпляр забирает повторно доставленное сообщение, он снова выполняет запись.
Доставку ровно один раз в распределённой системе практически невозможно гарантировать. Даже брокеры, заявляющие о поддержке семантики однократной обработки, гарантируют только те операции, которые они непосредственно контролируют, например доставку сообщений потребителям или запись данных обратно в брокер. Они не могут гарантировать внешние побочные эффекты, возникающие в результате действий потребителей в других системах. Надежное решение заключается не в том, чтобы устранять дублирующую доставку. Это нужно, чтобы потребитель мирился с этим. Когда вы сочетаете доставку как минимум один раз с потребителем, который игнорирует дубликаты, можно добиться фактически однократной обработки.
Решение
Сделайте потребителя идемпотентным, ведя учёт обработанных сообщений и пропуская любое сообщение, которое уже встречалось ранее. Потребитель основывает это решение на стабильном идентификаторе, который сохраняется при повторной доставке, проверяет постоянное хранилище, чтобы определить, был ли этот идентификатор уже обработан, и либо обрабатывает сообщение, либо отбрасывает его как дубликат.
Ниже описан основной поток.
- Прочитайте сообщение и извлеките его ключ дедупликации.
- Проверьте хранилище дедупликации для этого ключа.
- Если ключ существует, считайте сообщение дубликатом. Подтвердите его и остановите, при необходимости возвращая ранее записанный результат.
- Если ключ не существует, обработайте сообщение и запишите ключ в одной атомарной операции, а затем подтвердите сообщение.
Выбор стабильного ключа дедупликации
Ключ должен однозначно и последовательно определять логическое сообщение в каждой повторной версии. Используйте идентификатор сообщения, назначенный отправителем, или ключ идемпотентности на уровне бизнес-операции, который идентифицирует конкретную логическую операцию, а не разделяемый контекст корреляции, который может передаваться в нескольких сообщениях. В Служебная шина Azure свойство MessageId используется для этой цели, поскольку оно однозначно идентифицирует сообщение и его полезную нагрузку. Не используйте CorrelationId в качестве ключа, так как он группирует связанные сообщения, например запрос и ответы на него. Для событий, которые соответствуют спецификации CloudEvents, сочетание source и id атрибутов однозначно идентифицирует событие и остается стабильным в разных экземплярах.
Не следует опираться на идентификаторы транспортного уровня, которые брокер генерирует заново при повторной доставке, или на значения, производные от попыток доставки, поскольку эти значения различаются у дубликатов и мешают их обнаружению. Кроме того, избегайте извлечения ключа из переменных полей, таких как метки времени получения.
Если несколько независимых потребителей обрабатывают один канал, например несколько подписчиков в дизайне публикации и подписки, каждый потребитель законно обрабатывает собственную копию сообщения и должен независимо отслеживать завершение обработки сообщений. Если эти потребители совместно используют одно хранилище дедупликации, используйте для записи составной ключ из идентификатора потребителя и идентификатора сообщения. Хранилище, использующее только идентификатор сообщения в качестве ключа, позволяет первому потребителю подавить обработку сообщения всеми остальными.
Выбор места хранения обработанных ключей
У вас есть два распространенных варианта:
Выделенная таблица дедупликации. Потребитель поддерживает отдельную таблицу, иногда называемую inbox, в которой хранится по одной строке для каждого обработанного ключа. Такой подход позволяет отделять дедупликацию от бизнес-данных и хорошо работать, когда многие типы сообщений совместно используют один механизм.
Само юридическое лицо. Потребитель сохраняет ключ в записи, которую сообщение создает или обновляет. Этот подход позволяет избежать создания отдельной таблицы, но жёстко привязывает дедупликацию к структуре бизнес-данных.
Зафиксируйте маркер и побочные эффекты атомарно
В потоке «сначала проверка, затем обработка» есть окно отказа. Если потребитель обрабатывает сообщение, а затем на отдельном этапе записывает ключ, сбой между этими двумя операциями приводит к тому, что побочные эффекты уже применены, а ключ ещё не записан, поэтому при следующей доставке сообщение обрабатывается повторно.
Устраните это окно сбоя, записывая маркер дедупликации и побочные эффекты бизнес-логики в рамках одной и той же транзакции. Если обе операции фиксируются вместе или не фиксируются вовсе, то при повторной доставке либо обнаруживается зафиксированный маркер и обработка пропускается, либо маркер не обнаруживается, потому что транзакция была откатана, и сообщение можно безопасно обработать повторно. Этот транзакционный вариант представляет собой паттерн входящих сообщений и является парным паттерном на стороне потребления для паттерна транзакционных исходящих сообщений на стороне отправки.
Защита от одновременных дубликатов
При доставке как минимум один раз с несколькими конкурирующими потребителями два экземпляра могут одновременно получить копии одного и того же сообщения. Оба могут пройти проверку на существование прежде, чем любой из них выполнит фиксацию, поэтому сама по себе эта проверка не предотвращает двойную обработку.
Обеспечивайте корректность на уровне хранилища данных, а не в логике приложения:
Используйте уникальное ограничение для ключа дедупликации. Обе транзакции пытаются вставить ключ, но только одной это удаётся. Другой нарушает ограничение и обрабатывает сообщение как дубликат. Такой подход делает базу данных единственным арбитром в разрешении состояния гонки.
Избегайте состояний гонки при схеме «проверка с последующей установкой» в кэшах. Шаблон, который проверяет ключ, а затем задает его в двух отдельных операциях с окном, которое позволяет одновременно повторять оба утверждения ключа. Используйте атомарную условную запись, например вставку, которая завершается ошибкой при конфликте, или операцию set-if-absent, чтобы захват ключа выполнялся за один атомарный шаг.
Обработка побочных эффектов, которые не могут присоединиться к транзакции
Некоторые процессы не могут участвовать в транзакции базы данных потребителя, например вызов стороннего API или запись во внешнее хранилище. Для этих процессов используйте двухэтапный подход:
- Запишите ключ в состоянии в процессе перед тем, как выполнить внешнее действие.
- Выполните эту процедуру.
- Обновите запись, чтобы завершить и сохранить результат.
При повторной доставке завершённая запись позволяет пропустить повторное выполнение вызова. Незавершённая запись указывает на то, что предыдущая попытка могла быть выполнена лишь частично или всё ещё обрабатывается другим потребителем.
Проблемы и рекомендации
При принятии решения о том, как реализовать этот шаблон, учитывайте следующие моменты:
Предпочитайте изначально идемпотентные операции. Некоторые операции по своей природе идемпотентны и не требуют ведения учёта для дедупликации. Операция upsert, выполняемая по бизнес-идентификатору, запись, устанавливающая абсолютное значение, а не приращение, или HTTP-запрос
PUTк идентификатору ресурса — любая из этих операций дает один и тот же результат независимо от того, выполняется ли она один раз или многократно.Иногда операцию можно сделать идемпотентной по своей природе с помощью передачи состояния в событии, когда сообщение содержит итоговое абсолютное состояние, например новый статус заказа, так что потребитель применяет его как операцию вставки или обновления, а не как относительное изменение.
Совет
Сначала закладывайте естественную идемпотентность в проектирование и добавляйте методы дедупликации только для операций, которые нельзя сделать естественно идемпотентными.
Управление жизненным циклом записей дедупликации. Записи дедупликации накапливаются, если вы не удаляете их по истечении срока хранения. Сохраняйте каждую запись как минимум в течение того времени, в течение которого брокер может повторно доставить исходное сообщение. Это окно зависит от максимального числа попыток доставки, заданного брокером, тайм-аута блокировки или видимости и времени жизни сообщения. Задайте для записей дедупликации время жизни (TTL), которое превышает это окно, чтобы поздняя повторная доставка всё ещё могла найти свой маркер. Удаление записей слишком рано открывает окно для дубликатов. Учитывайте сообщения, которые оператор повторно отправляет из очереди недоставленных сообщений, поскольку повторная отправка может произойти значительно позже обычного окна повторной доставки.
Используйте фреймворк обмена сообщениями вместо того, чтобы реализовывать дедупликацию самостоятельно. Корректная реализация хранилища дедупликации, атомарной фиксации и очистки записей чревата ошибками. Платформы на основе сообщений предоставляют этот шаблон как встроенную функцию.
Например, NServiceBus дедупликирует входящие сообщения по идентификатору сообщения и обеспечивает настраиваемое хранение и очистку для данных дедупликации. В папке "Входящие" потребитель MassTransit отслеживает полученные сообщения по идентификатору сообщения, чтобы обеспечить точное поведение потребителя.
Дедупликация на стороне брокера снижает, но не устраняет потребность в логике идемпотентной обработки на стороне потребителя. Некоторые платформы фильтруют дубликаты на уровне транспорта. Служебная шина Azure обнаружение дубликатов отбрасывает сообщения с повторяющимся
MessageIdв пределах настроенного временного окна, что предотвращает появление дубликатов, вызванных повторными попытками отправки со стороны отправителя. Эта функция работает на стороне отправителя и в пределах ограниченного окна. Это не предотвращает повторную обработку одного и того же сообщения потребителем после повторной доставки, поэтому вам по-прежнему нужна идемпотентная логика потребителя. Считайте функции платформы первым уровнем защиты, который снижает объём дубликатов, а не как замену паттерна.Учетная запись для упорядочивания сообщений. Дедупликация удаляет дубликаты, но не гарантирует порядок. Если для потребителя важен порядок обработки, используйте этот шаблон вместе с механизмом упорядочивания, например с сеансами сообщений Служебная шина Azure, или добавьте данные о последовательности или версии, чтобы потребитель мог отклонять устаревшие сообщения.
Инструмент наблюдаемости. Выведите ключ дедупликации и идентификатор корреляции в структурированных журналах и отслеживайте метрику для обнаруженных повторяющихся данных. Рост уровня дублирования может указывать на неправильную конфигурацию производителя, слишком маленькое окно подтверждения или блокировки либо на неисправных потребителей. Используйте сквозную трассировку и корреляцию, чтобы отслеживать сообщение при его прохождении через службы.
Передавайте идемпотентность в последующие вызовы. Если сделать одного потребителя идемпотентным, это не защитит службы, которые он вызывает. Когда потребитель в рамках обработки вызывает нижележащие сервисы, передавайте ключ идемпотентности, чтобы каждый уровень мог устранять дублирование в своей работе.
Когда следует использовать этот шаблон
Используйте этот шаблон, когда:
Вы обрабатываете сообщения от брокера, который обеспечивает доставку как минимум один раз, что является режимом по умолчанию для большинства брокеров.
При повторной обработке сообщения возникают неправильные результаты, такие как повторяющиеся финансовые транзакции, создание повторяющихся ресурсов или повторяющиеся уведомления.
Несколько конкурирующих потребителей обрабатывают один канал, что делает одновременную повторяющуюся доставку вероятной.
Этот шаблон может быть не подходит, если:
Каждая операция, выполняемая потребителем, уже по своей природе идемпотентна, поэтому повторная обработка безвредна, а ведение учёта для дедупликации влечёт дополнительные издержки, не принося никакой пользы.
Рабочая нагрузка допускает последствия эпизодической повторной обработки, а стоимость хранилища для дедупликации превышает ущерб от появления дубликата.
Идемпотентная обработка за рамками обмена сообщениями
Этот шаблон применяет идемпотентность к потребителям сообщений, но идемпотентная обработка является более широким принципом надежности. Любая операция, которую можно выполнять более одного раза над одной и той же задачей, только выигрывает от этого. Этот принцип включает ETL-преобразования (извлечение, преобразование и загрузка данных), которые повторно обрабатывают данные при повторном воспроизведении, потоковую обработку, возобновляемую с контрольной точки, запланированные задания, которые перекрываются или перезапускаются, а также конечные точки webhook или HTTP, получающие повторные доставки.
В каждом случае применяется один и тот же основной метод:
- Идентифицируйте единицу работы со стабильным ключом.
- Запишите то, что вы уже обработали.
- Пропустить или поглощать дубликаты, чтобы повторение работы не изменило результат.
Механизмы, приведенные в этой статье, такие как стабильные ключи, атомарные маркеры и уникальные ограничения, передаются в эти контексты даже при отсутствии посредника сообщений.
Проектирование рабочей нагрузки
Оцените, как использовать шаблон Идемпотентного потребителя в проектировании рабочей нагрузки для решения целей и принципов, описанных в основных принципах Azure Well-Architected Framework. В следующей таблице приведены рекомендации по использованию этого шаблона для целей каждого компонента.
| Столп | Как этот шаблон поддерживает цели основных компонентов |
|---|---|
| Решения по проектированию надежности помогают рабочей нагрузке стать устойчивой к сбоям и гарантировать, что она восстанавливается до полнофункционального состояния после сбоя. | Этот шаблон позволяет рабочей нагрузке использовать по крайней мере один раз доставку и безопасные повторные попытки без повреждения данных, которая превращает повторяющуюся доставку из риска правильности в допустимое условие. - RE:07 Самосохранение - Обработка временных сбоев |
Если этот шаблон вводит компромиссы внутри столпа, рассмотрите их против целей других столпов.
Пример
В следующем примере приведен идемпотентный потребитель, который обрабатывает заказы из Служебная шина Azure и сохраняет состояние в Azure Cosmos DB for NoSQL.
Производитель устанавливает для служебная шина MessageId идентификатор заказа бизнес-уровня. Потребитель получает сообщения в режиме PeekLock , который повторно удаляет сообщение, если потребитель не завершает его в течение длительности блокировки. Контейнер Azure Cosmos DB потребителя разбивается на разделы по идентификатору заказа (/orderId), а документу id присваивается тот же идентификатор заказа, поэтому каждая копия данного заказа попадает в один и тот же логический раздел, а сама запись заказа служит маркером дедупликации.
Потребитель обрабатывает каждое сообщение следующим образом:
- Прочтите сообщение и используйте его
MessageIdв качестве ключа дедупликации. - Создайте документ заказа, установив для
idи ключа секции значение идентификатора заказа. - Если создание прошло успешно, завершите обработку сообщения, чтобы служебная шина удалил его из очереди.
- Если создание завершается ошибкой со статусом HTTP 409 (Конфликт), потому что документ с таким
idуже существует, прочитайте существующий документ и сравните его с текущим сообщением. Если сохранённый хэш запроса или неизменяемые бизнес-поля совпадают, считайте сообщение дубликатом, завершите его обработку и пропустите дальнейшую обработку. Если они не совпадают, производитель, возможно, повторно использовал идентификатор для другого содержимого, либо сведения о заказе могли измениться с момента его первой обработки, поэтому отправьте сообщение в очередь недоставленных сообщений или сгенерируйте оповещение, а не молча отбрасывайте его. - Если обработку не удаётся выполнить из-за временного сбоя, откажитесь от сообщения, чтобы служебная шина повторно доставил его, или дождитесь истечения срока блокировки, чтобы его получил другой потребитель.
Операция создания атомарна, поэтому она служит и проверкой на наличие дубликатов, и операцией записи. Два потребителя, получающие копии одного сообщения, не могут создать заказ. Один создает победу, а другой возвращает конфликт и безопасно удаляет его дубликаты.
Если в ходе обработки необходимо записать более одного документа, используйте транзакционный пакет, включающий как документ дедупликации, так и бизнес-документы с одним и тем же ключом раздела. Так как транзакционный пакет работает в пределах одного логического раздела, выберите ключ раздела, общий для всех документов одного сообщения. Пакет либо фиксирует все документы сразу, либо не фиксирует ни одного, поэтому сбой между обработкой и подтверждением не может привести к тому, что маркер дедупликации и бизнес-данные окажутся в рассинхронизированном состоянии. Если пакет пытается создать документ, который уже существует, возвращается код состояния 409 (Conflict), что позволяет выявить дубликат.
Чтобы сделать этот потребитель устойчивым к повторяющимся повторным попыткам отправки, включите обнаружение повторяющихся данных в очереди. Обнаружение дубликатов предотвращает повторную отправку сообщений в пределах своего окна истории, а идемпотентный потребитель обрабатывает все дубликаты, которые выходят за пределы этого окна или возникают в результате повторной доставки.
Следующий шаг
- Параметры асинхронного обмена сообщениями в Azure описывают варианты инфраструктуры обмена сообщениями, определяющие гарантии доставки и требования к повторяющейся обработке.
Связанные ресурсы
Паттерн Transactional Outbox представляет собой сторону издателя в этом шаблоне. Он надежно публикует сообщения, фиксируя их в той же транзакции, что и бизнес-данные.
Шаблон повторных попыток позволяет приложениям обрабатывать временные ошибки путем повторных операций, что делает идемпотентную обработку необходимой, так как повторные попытки могут привести к дублированию доставки.
Отказоустойчивая архитектура Центры событий Azure и Функции Azure применяет этот шаблон к функциям, запускаемым Центры событий Azure, включая методы дедупликации в потоках событий.
Проектирование Функции Azure для идентичных входных данных предоставляет рекомендации по созданию идемпотентных функций, которые допускают повторяющиеся вызовы.