Последовательный паттерн конвоя

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

Этот шаблон устраняет напряженность между поддержанием правильности первого входа (FIFO) в каждой логической группе и масштабированием параллельной обработки между группами. Проект гарантирует, что ограничения упорядочивания не становятся узким местом на уровне системы.

Контекст и проблема

Приложения часто должны обрабатывать связанные сообщения в том порядке, в который они приходят, при этом масштабирование по-прежнему выполняется для обработки повышенной нагрузки. В распределённой архитектуре это требование трудно выполнить, поскольку рабочие процессы независимо извлекают сообщения из общей очереди. Когда несколько обработчиков конкурируют за получение сообщений, как в паттерне Competing Consumers, порядок обработки нарушается.

Рассмотрим систему отслеживания заказов, которая получает поток операций, например создание заказа, добавление транзакции, изменение прошлой транзакции и удаление заказа. Операции каждого заказа должны обрабатываться в порядке FIFO, так как применение их вне последовательности приведет к повреждению состояния заказа. Однако во входящей очереди операции по множеству заказов перемежаются. Один потребитель, обеспечивающий глобальный порядок обработки, становится узким местом, а несколько потребителей могут обрабатывать операции одного и того же заказа не по порядку.

Каждый из простых подходов к этой проблеме по-своему не срабатывает:

  • Один потребитель. Один потребитель сохраняет порядок сообщений, так как обрабатывает одно сообщение одновременно, но не может масштабироваться для обработки повышенной пропускной способности.

  • Несколько конкурирующих потребителей. Несколько потребителей масштабируют пропускную способность путем параллельного извлечения сообщений, но они теряют гарантии заказа на группу. Два обработчика могут получать последовательные сообщения для одного и того же заказа и обрабатывать их одновременно или не по порядку, что нарушает состояние заказа.

Решение

Последовательный шаблон конвоя секционирует связанные сообщения в категории и обрабатывает каждую категорию последовательно, по одному сообщению одновременно, а категории обрабатываются параллельно.

Шаблон работает путем назначения каждого сообщения ключа категории, определяющего группу, к которой она принадлежит. Брокер сообщений использует этот ключ для секционирования сообщений в логические группы. В каждой группе брокер применяет упорядочение FIFO, чтобы потребитель, который блокирует группу, получает сообщения строго в последовательности, в которую они были вложены. Разные группы могут одновременно обрабатываться разными потребителями, поэтому система горизонтально масштабируется между группами без потери порядка внутри каждой отдельной группы.

В службе Служебная шина Azure сеансы сообщений предоставляют встроенную реализацию этого шаблона.

На следующей схеме показан общий шаблон последовательного конвоя.

Схема последовательного конвоя. В нем показан производитель, центральная очередь и три потребителя.

Схема проходит слева направо и показывает три компонента. Слева находится блок с надписью «producer». Стрелка направлена от производителя к центральному блоку с меткой «queue». От очереди вправо указывают три стрелки, каждая помечена названием категории. Стрелка вверху помечена категорией A и указывает на поле, помеченное потребителем A в правом верхнем углу. Средняя стрелка имеет метку «категория B» и указывает на блок с меткой «потребитель B» в центре справа. Стрелка внизу помечена категорией C и указывает на поле, помеченное потребителем C в правом нижнем углу. Три отдельные линии от очереди к трем потребителям иллюстрируют, что очередь разделяет сообщения по ключу категории. Он направляет сообщения каждой категории исключительно соответствующему потребителю, что позволяет каждому потребителю обрабатывать соответствующие потоки сообщений параллельно и в порядке FIFO без вмешательства друг друга.

В очереди сообщения для разных категорий могут чередоваться, как показано на следующей схеме.

Схема с четырьмя категориями чередуемых сообщений в одной очереди. Каждая категория занимает собственную горизонтальную полосу.

На схеме показан интерьер одной очереди, отображаемой в виде широкой прямоугольной области. Внутри очереди четыре горизонтальные полосы стекаются по вертикали и нумеруются 1–4 по правому краю, каждая со стрелкой, указывающей направо, чтобы указать направление потока сообщений. Lane 1 содержит четыре блока сообщений, распределенные по всей ширине очереди, что указывает на большой объем сообщений для этой категории. Полоса 2 также содержит четыре блока сообщений, расположенных по всей ширине. Lane 3 содержит три блока сообщений, а полоса 4 содержит один блок сообщения, расположенный в правой части очереди. Различные положения и плотность блоков сообщений в разных дорожках показывают, что сообщения из всех четырёх категорий поступают вперемежку в одну общую очередь. Несмотря на это чередование, сообщения каждой категории сохраняют порядок поступления слева направо в пределах своей полосы, что показывает, что порядок FIFO для каждой категории сохраняется, даже когда разные категории разделяют единую структуру очереди.

Этот шаблон предоставляет несколько ключевых преимуществ:

  • Последовательная обработка по группам. Сообщения в каждой категории обрабатываются строго в последовательности, что предотвращает условия гонки, изменения состояния вне порядка и необходимость переупорядочения обходных решений.

  • Горизонтальное масштабирование между группами. Каждая категория является независимой единицей параллелизма. Добавление потребителей увеличивает пропускную способность пропорционально количеству активных категорий без нарушения гарантий упорядочения.

  • Развязка производителя и потребителя. Производители помещают сообщения в очередь, не зная, какой потребитель будет их обрабатывать и когда. Потребители являются независимо масштабируемыми и заменяемыми.

Проблемы и рекомендации

Рассмотрим следующие моменты, когда вы решите, как реализовать этот шаблон:

  • Категория и единица масштабирования. Определите, по какому свойству входящих сообщений можно выполнять горизонтальное масштабирование. Ключ категории определяет единицу параллелизма: каждое отдельное значение ключа становится независимо обрабатываемой группой. В сценарии отслеживания заказов это свойство является идентификатором заказа. Выбор слишком грубого ключа (например, одного идентификатора клиента для всех заказов) ограничивает параллелизм, тогда как выбор слишком детального ключа не даёт существенного преимущества с точки зрения упорядочивания.

  • Ограничения пропускной способности. Оцените пропускную способность целевого сообщения. Так как этот шаблон обеспечивает последовательную обработку в каждой категории, пропускная способность для каждой категории ограничивается временем обработки одного сообщения. Оптимизируйте время обработки каждого сообщения, например с помощью асинхронного ввода-вывода или пакетирования записей в нижестоящие системы, поскольку именно это время напрямую определяет максимальную пропускную способность для каждой категории. Если общая требуемая пропускная способность очень высока, пересмотрите, действительно ли необходим строгий порядок FIFO на протяжении всего жизненного цикла сообщений. Альтернативы включают обязательное использование начального и конечного сообщений для обозначения границ последовательности или сортировку сообщений по меткам времени в пределах окна пакетной обработки с последующей отправкой пакета на параллельную обработку.

  • Возможности службы. Убедитесь, что ваш выбор брокера сообщений поддерживает однократную обработку сообщений в очереди или категории очереди. Не все службы обмена сообщениями обеспечивают блокировку на уровне сессии или гарантии FIFO в пределах партиции. Если брокер не поддерживает эту возможность встроенными средствами, потребитель должен реализовать собственную логику координации, что повышает сложность и создаёт риск дублирования обработки, пропуска сообщений или выполнения не по порядку. Поддержка сеансов также может ограничить выбор уровня обмена сообщениями или номера SKU, что влияет на стоимость.

  • Эволюционируемость. Планирование добавления новых категорий сообщений в систему. Шаблон должен соответствовать росту кратности категорий, не требуя структурных изменений потребителям. Например, предположим, что система реестра, описанная ранее, относится к одному клиенту. Если вам нужно подключить нового клиента, вы сможете добавить набор процессоров реестра, которые распределяют работу на идентификатор клиента, не изменяя топологию очередей.

  • Доставка сообщений вне заказа. Сообщения могут приходить не по порядку из-за переменной сетевой задержки между производителем и брокером до того, как начнёт действовать упорядочивание в рамках сеанса на стороне брокера. Рассмотрите возможность использования порядковых номеров для проверки последовательности в каждой категории. Вы также можете включить флаг завершения последовательности в последнее сообщение транзакции, чтобы потребители могли определить, когда последовательность завершена.

  • Обработка подозрительных сообщений. Сообщение, которое неоднократно завершает обработку в сеансе, блокирует все последующие сообщения в этом сеансе, так как шаблон применяет строгое последовательное упорядочение. Разработайте стратегию обнаружения проблемных сообщений, например, отслеживая количество попыток доставки, и перемещайте их в очередь недоставленных сообщений после достижения заданного порога повторных попыток, чтобы можно было продолжить обработку остальных сообщений в сеансе.

  • Доступность брокера. Брокер сообщений — это общая зависимость для всех категорий. Его доступность и устойчивость непосредственно влияют на гарантии надежности шаблона. Оцените функции отказоустойчивости на уровне брокера, такие как зоны доступности и геораспределённое аварийное восстановление, исходя из требований к доступности рабочей нагрузки и бюджетных ограничений, поскольку конфигурации с более высокой отказоустойчивостью обычно увеличивают затраты.

  • Корректность ключа производителя. В шаблоне предполагается, что производители правильно устанавливают ключ категории (идентификатор сеанса) для каждого сообщения. Если производитель задает неправильный ключ либо случайно, либо из-за ошибки, сообщение направляется в неправильный сеанс и повреждает состояние этой группы. Убедитесь, что производители последовательно назначают ключи категорий, и рассмотрите возможность добавить на стороне потребителя логику проверки ключей, если последствия неверно маршрутизированного сообщения могут быть серьёзными.

  • Операционная сложность. Мониторинг обработки с использованием сеансов влечёт дополнительные операционные накладные расходы по сравнению с обычной обработкой очереди. Операторам необходимо иметь представление о накоплении сообщений в сеансах (количестве активных сеансов и количестве сообщений, ожидающих обработки в каждом сеансе), чтобы выявлять категории, которые отстают. Для недоставленных сеансов требуется отдельный рабочий процесс мониторинга и исправления, чтобы исследовать неудачные сообщения, устранять первопричину и воспроизводить исправленные сообщения обратно в сеанс.

  • Состязание за блокировку сеанса и задержка. Блокировка сеанса приводит к дополнительным временным затратам, так как каждый потребитель перед обработкой сообщений должен получить эксклюзивную блокировку для сеанса. Когда потребитель держит блокировку сеанса, ни один другой потребитель не может обрабатывать сообщения из этого сеанса, даже если потребитель медленно или временно застопорился. Если длительность блокировки слишком мала, истечение срока блокировки может привести к повторной обработке сообщения. Если срок блокировки слишком велик, зависший потребитель задерживает восстановление. Настройте длительность блокировки сеанса на основе ожидаемого времени обработки сообщений и реализуйте продление блокировки для длительных операций.

  • Потребительское масштабирование и стоимость. Параллелизм в рамках сеансов означает одновременную работу экземпляров потребителей. В бессерверной модели, такой как Функции Azure, каждый активный сеанс сопоставляется с одновременным выполнением, а в выделенной модели он сопоставляется с экземпляром или потоком. Таким образом, количество активных сеансов напрямую влияет на затраты на вычисления. Планируйте ограничения масштабирования потребителей и средства управления параллелизмом, чтобы сбалансировать пропускную способность и затраты.

Когда следует использовать этот шаблон

Используйте этот шаблон, когда:

  • Сообщения поступают в определённом порядке и должны обрабатываться в том же порядке.
  • Сообщения можно классифицировать таким образом, чтобы каждая категория стала независимой единицей масштабирования для системы.

Этот шаблон может быть не подходит, если:

  • Вы ожидаете чрезвычайно высокую пропускную способность (миллионы сообщений в минуту), так как требование FIFO ограничивает масштабирование, которое может достичь система.

  • Порядок сообщений не обязателен. Когда сообщения можно обрабатывать независимо в любом порядке, шаблон конкурирующих потребителей обеспечивает более простое горизонтальное масштабирование без координации нагрузки на блокировку сеансов.

Проектирование рабочей нагрузки

Оцените, как использовать последовательный конвой в проектировании рабочей нагрузки для решения целей и принципов, описанных в основных принципах Azure Well-Architected Framework. В следующей таблице приведены рекомендации по использованию этого шаблона для целей каждого компонента.

Столп Как этот шаблон поддерживает цели основных компонентов
Решения по проектированию надежности помогают рабочей нагрузке стать устойчивой к сбоям и гарантировать, что она восстанавливается до полнофункционального состояния после сбоя. В этом шаблоне используется FIFO-упорядочение на основе сеансов, чтобы устранить условия гонки, логику обработки сообщений, в которой часто возникают конфликты, и другие обходные решения для обработки сообщений, поступающих в неправильном порядке, которые могут приводить к сбоям в работе.

- Критически важные потоки RE:02
- RE:07 Фоновые задания

Если этот шаблон вводит компромиссы внутри столпа, рассмотрите их против целей других столпов.

Example

В Azure этот шаблон можно реализовать с помощью сеансов сообщений служебная шина. Для получателей можно использовать либо Azure Logic Apps с соединителем служебная шина в режиме peek-lock, либо Функции Azure с триггером служебная шина.

Когда производитель задает для сообщения свойство SessionId, служебная шина группирует все сообщения, которые имеют один и тот же идентификатор сеанса, в один логический сеанс. Потребитель принимает сеанс и получает исключительную блокировку на него. Эта блокировка гарантирует, что в каждый момент времени сообщения в этом сеансе обрабатывает только один потребитель и что сообщения поступают в порядке FIFO. Другие потребители могут одновременно принимать и обрабатывать различные сеансы, обеспечивая параллельную пропускную способность между группами.

В примере отслеживания заказов система обрабатывает каждое сообщение реестра в том порядке, в котором оно получено, и отправляет каждую транзакцию в другую очередь, в которой для категории задан идентификатор заказа. В этом сценарии транзакция никогда не распространяется на несколько заказов, поэтому консьюмеры обрабатывают категории параллельно, но внутри каждой категории — в порядке FIFO.

Обработчик реестра распределяет сообщения, разделяя пакетное содержимое каждого сообщения в первой очереди:

Схема архитектуры последовательного конвоя. В нем показан производитель, очередь реестра, обработчик реестра, очередь транзакций и три обработчика заказов.

Схема состоит из пяти этапов и читается слева направо. Слева, у самого края, находится блок с надписью «producer». Стрелка ведет от производителя к блоку с надписью «ledger queue». Стрелка ведет от очереди журнала к блоку с надписью «обработчик журнала». Стрелка направлена от обработчика реестра к блоку с надписью «очередь транзакций». От очереди транзакций вправо направлены три стрелки, каждая помечена категорией сеанса. Верхняя стрелка помечена как «транзакции заказа A» и указывает на блок с меткой «обработчик заказа A». Средняя стрелка помечена как «транзакции заказа B» и указывает на блок с меткой «обработчик заказа B». Нижняя стрелка помечена как «транзакции заказа C» и указывает на блок с меткой «обработчик заказа C». Схема иллюстрирует переход от последовательной обработки к параллельной: производитель последовательно отправляет все сообщения через очередь реестра и обработчик реестра, а обработчик реестра присваивает идентификатору сеанса каждого сообщения соответствующий идентификатор заказа, прежде чем поместить сообщение в очередь транзакций. Затем очередь транзакций направляет сообщения каждого заказа исключительно в соответствующий обработчик заказов, который позволяет обработчику заказов A, обработчику заказов B и обработчику заказов C использовать соответствующие сеансы параллельно и в порядке FIFO.

Обработчик реестра выполняет три шага:

  1. Проходит по реестру по одной транзакции за раз.
  2. Задает идентификатор сеанса сообщения в соответствии с идентификатором заказа.
  3. Отправляет каждую транзакцию реестра в вторичную очередь с идентификатором сеанса, заданным для идентификатора заказа.

Потребители отслеживают вторичную очередь и обрабатывают все сообщения с совпадающими идентификаторами заказов в порядке FIFO. Пользователи используют режим peek-lock.

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

Поддержка технологий

  • Сеансы сообщений служебная шина: группирует сообщения по идентификатору сеанса и обеспечивает обработку в порядке FIFO в рамках каждого сеанса. Сеансы сообщений — это основной механизм Azure для реализации шаблона последовательного конвоя.

  • триггер Функции Azure служебная шина: поддерживает триггеры на основе сеансов, которые позволяют экземплярам функций обрабатывать сообщения из одного сеанса одновременно.

  • Соединитель Logic Apps для служебная шина: предоставляет соединитель служебная шина с поддержкой режима peek-lock для обработки сообщений из очередей с поддержкой сеансов в обработке на основе рабочих процессов.

Соавторы

Корпорация Майкрософт поддерживает эту статью. Следующие авторы написали эту статью.

Основной автор:

  • Нага Венката Черуву | Старший архитектор облачных решений и инфраструктура искусственного интеллекта

Чтобы увидеть непубличные профили в LinkedIn, войдите в LinkedIn.

  • Шаблон конкурирующих потребителей: несколько потребителей извлекают сообщения из общей очереди параллельно, что повышает пропускную способность, но удаляет гарантии порядка сообщений. Паттерн Sequential Convoy устраняет проблему с порядком обработки, которую вносит паттерн Competing Consumers. Он устраняет этот разрыв путем секционирования сообщений в сеансы с ключами категорий и последовательной обработки каждого сеанса.

  • Шаблон выравнивания нагрузки на основе очереди: очередь служит буфером для задач между производителями и потребителями, чтобы поглощать всплески нагрузки и сглаживать её неравномерность. Шаблон Sequential Convoy основан на этом механизме буферизации и дополняет его разделением на основе сеансов, благодаря чему очередь одновременно выравнивает нагрузку между категориями и сохраняет порядок FIFO внутри каждой категории.

  • Шаблон очереди приоритета: сообщения направляются в отдельные очереди или имеют приоритет в очереди, чтобы работа с более высоким приоритетом обрабатывалась до работы с более низким приоритетом. Если также необходимо сохранить порядок обработки внутри уровня приоритета, шаблон Sequential Convoy можно комбинировать с очередью с приоритетами, чтобы обеспечить обработку в порядке FIFO в каждом сеансе, определяемом ключом приоритета.

  • Peek-Lock сообщение (неразрушительное чтение): эта операция атомарно извлекает и блокирует сообщение из очереди или подписки для обработки.

  • Упорядоченная доставка коррелированных сообщений в Logic Apps с использованием сеансов служебная шина: в этой записи блога описывается поддержка шаблона Sequential Convoy в Logic Apps.