Герой нашего сегодняшнего рассказа — скромный и, на первый взгляд, заурядный диспетчер1. В наше время таких на заграничный манер все чаще назовут «роутер», поэтому, надеюсь, читатель простит меня за то, что буду использовать оба термина. Обосновался наш диспетчер на самой верхушке каталога процессов в Camunda-кластере сервиса бронирования жилья. Единственная его трудовая обязанность за целую смену — разбирать входящую корреспонденцию и направлять конечным обработчикам: работа не пыльная и все больше однообразная. Примечательного в нем на первый взгляд мало: аккуратный, документированный и до педантизма лаконичный, в таком и сломаться-то вроде нечему. Впрочем, смотрите сами2:

Протагонист
Протагонист

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

Как устроен роутер

Кратко про Camunda для тех, кто не в танке

Camunda 8 — движок бизнес-процессов, ядро которого — Zeebe — было для восьмерки написано с нуля. Поэтому восьмая версия — это не апгрейд седьмой, а полностью другой инструмент. Многие изыскания и аргументы из этой статьи к седьмой версии применимы быть не могут ввиду радикально разной архитектуры.

Начиная с версии 8.8 отдельные UI web-приложения: Operate (мониторинг и отладка) и Tasklist (ручной запуск процессов, управление задачами, веб-формы) — перекочевали в единый Orchestration Cluster.

Вся логика работы описывается одним предложением: забирать из Kafka-топика обращения (тикеты) пользователей, зарегистрированные первой линией поддержки, читать поле «тип обращения» и передавать процессу, обрабатывающему этот тип.

Один в поле

Главное, что нужно сказать о нашем герое, — он одиночка, в любой момент времени функционирует только один экземпляр (инстанс) процесса. Завершая работу, роутер отправляет сигнал своему преемнику, который незамедлительно встает на смену. На случай сбоя или первого включения оператор может запустить роутер вручную.

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

Передача смены
Передача смены

Вся эта конструкция держится на строгом порядке событий в BPMN: kill-сигнал транслируется до того, как экземпляр-эмитент сам дошел до этапа подписки на этот сигнал. Zeebe сигналы не буферизует — его получают только те подписки, которые открыты в момент рассылки, проверка «а не убью ли я себя сам?» не нужна.

Неопасный баг конструкции

Поскольку наш механизм подразумевает частичное перекрывание, а не полное взаимоисключение, существует ненулевая вероятность одновременного запуска двух роутеров: оба отправят kill-сигнал до того, как подпишутся на него, и выживут. Такой риск приемлем, потому что вероятность невероятно мала, а последствия безопасны — экземпляры будут просто разбирать топик вдвоем. Лечится же штатным механизмом ручной остановки.

«Сиди и слушай»

Убедившись, что никого вокруг не осталось, роутер приземляется на основное рабочее место, роль которого выполняет пользовательская задача (user task — BPMN-элемент, предназначенный для выполнения человеком). Вся центральная конструкция процесса читается как «сиди смирно и слушай».

Казалось бы, последнее, что мы хотим от роутера, — это чтоб он раздавал задачи людям, а последнее, чего мы хотим от людей, — чтоб они вмешивались в процесс. Но в этом и кроется одна из неочевидных уловок нашего героя — дизайн его схемы как раз полагается на то, что пользовательскую задачу никто в нормальном режиме не выполнит («сидение смирно»), одновременно давая человеку точку вмешательства на случай, если что-то пошло не по плану.

Но все же основное штатное назначение задачи — служить посадочной площадкой для пограничных событий, уводящих роутера со смены. А то самое «и слушай» реализует non-interrupting message event subprocess — как бы пугающе он ни звучал, событийный подпроцесс — это хрестоматийный способ канонично смоделировать «обработать много событий за окно».

Одно окно и три двери

Прерывающее событие (interrupting, сплошная линия) фактически перенаправляет процесс: гасит текущий поток и начинает собственный. Непрерывающее (non-interrupting, пунктирная линия) на основную работу не влияет — сидит в фоне и, пока активно, принимает все адресованные ему сообщения, на каждое порождая отдельный параллельный поток.

В нашем роутере при деле оба типа, но расселены они по разным местам. На границе пользовательской задачи (boundary events) висят два прерывающих события — таймер и аварийный сигнал; вместе с ручным «выполнением» самой задачи они и образуют «три двери» на выход. А «слушает» за роутера непрерывающее стартовое событие событийного подпроцесса: оно подписано на Kafka всё время жизни экземпляра и на каждое входящее сообщение поднимает свой подпоток.

Посадочная площадка
Посадочная площадка

Благодаря непрерывающему Kafka-коннектору, наш роутер работает как централизованный коллектор. Сама же маршрутизация полученного сообщения тривиальна и сводится к else-if шлюзу с дефолтным маршрутом к сборщику мусора (защищает от застрявших исполнений).

Завершается же смена роутера одним из трех способов:

  1. Штатный режим — таймер. TTL экземпляра вшит в определение процесса, но может быть переопределен при ручном запуске. В нашем случае 24-часовая вахта была оптимальной.

  2. Ручной режим — основной поток исполнения. Оператор «выполняет» назначенную задачу «Остановить роутер».

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

Заметьте, что в последнем случае процесс просто завершается (т.н. none end event) — нас остановил другой уже активный процесс, поэтому передавать смену некому (можно сказать, что ее перехватили). Ручной и штатный режим по умолчанию всегда запускают новый экземпляр, none завершение в этом случае возможно, только если оператор заранее определил такое поведение — работающих роутеров в этом случае не остается, и сообщения будет некому разбирать.

Вход и выход

Форму любого процесса определяют требования среды, в которой он выполняется:

  1. Все тикеты пользователей публикуются в общий топик Kafka — кто-то обязан их распределить.

  2. Cost-per-instance — любой экземпляр, не делающий бизнес-работы, — это накладные расходы, оплаченные деньгами.

  3. Конечные обработчики тикетов живут часами и неделями — и каждый должен быть самостоятельным наблюдаемым экземпляром.

  4. Роутер должен быть управляем — вмешательство человека должно быть частью штатного функционирования, без редеплоя.

  5. Пиковая нагрузка 10 RPS — разрешает принять ограничения и не гнаться за избыточным масштабированием.

Для таких сугубо интеграционных компонентов при проектировании ключевыми осями решений являются: как поймать сообщение (вход) и как передать его дальше (выход).

Как поймать сообщение

Если опустить вариант с ручным запуском процесса, остается три принципиальных способа реализации:

  1. Без роутера

  2. Долгоживущий роутер

  3. Короткоживущий роутер

Как читать топик?
Как читать топик?

Маршрутизация без роутера

Camunda позволяет сделать Kafka-коннектор стартовым событием всего процесса, с фильтрацией сообщений по условию активации. Каждый обработчик знает, какие тикеты его (ticket.type=='someType'), роутер вообще не нужен! Круто же?

Чтобы фильтрация работала, каждый обработчик должен видеть каждое сообщение. Это значит, что вместо одной консьюмер-группы роутера нам нужны отдельные на каждый обработчик. N обработчиков — N чтений топика, то есть вместо O(M) мы получаем O(N x M). И каждое сообщение будет десериализовано и проверено через FEEL3, только чтоб пустить в обработку лишь небольшую порцию всего потока.

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

Синглтон роутер

Синглтон — не самый распространенный в BPM паттерн: в Zeebe нет механизма взаимного исключения4, нельзя напрямую, не выходя наружу, спросить «а не запущен ли уже такой экземпляр?». Инвариант «один роутер за раз» реализуем только самой его формой, и единственный локальный инструмент, который у нас для этого есть, — сигнал, BPMN-элемент, оптимальный для fire-and-forget.

Причем это не обходной путь вместо блокировки — это лучше. Блокировка, удерживаемая умершим роутером, требует lease и expiry, прежде чем кто-то другой сможет продолжить. Здесь нет ни держателя, ни того, что надо освобождать: новейший экземпляр побеждает по построению — и это правильная политика для компонента, не несущего собственного состояния.

Зачем вообще периодически пересоздавать экземпляр? Почему бы не взять диспетчера в кабалу и пусть работает без остановки?

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

Регулярный перезапуск «сбрасывает» накопленное, завершенные процессы движок переносит в архив. А в случае локальных сбоев handover работает еще и как механизм самолечения. Нам лишь надо обеспечить безопасное окно передачи. Для этого при настройке Kafka-коннектора нужно установить Message TTL — время, в течение которого сообщение буферизуется в Zeebe в ожидании корреляции с экземпляром.

Router-per-message

Но к чему эта свистопляска с синглтоном, если роутер можно стартовать с Kafka-коннектора на каждом сообщении?

Вот тут в дело вступает экономика и лицензионная политика Camunda: биллинговая метрика «Root Process Instance» начисляет копеечку за каждый запущенный корневой экземпляр процесса. С линейными роутерами обработка каждого сообщения будет требовать два экземпляра: роутер + обработчик — стоимость решения вырастает вдвое!

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

Замечания про биллинг
  • Двойного биллинга в router-per-message можно избежать, используя Call Activity, о котором речь пойдет дальше, поэтому этот аргумент не бескомпромиссный

  • На стоимость синглтона влияет не только handover TTL, который меняет частоту создания экземпляров, но и конфигурация пользовательской задачи, которая может вовлечь другую биллинговую метрику — «Task User»

Как передать сообщение

В нашем распоряжении 4 способа организовать делегирование логики:

  1. Sub-process. Не является отдельным процессом, он вложен в родительский; это скорее аналог блоков кода вроде for-циклов.

  2. Call activity. Самый очевидный go-to вариант для такого рода задач, аналогичен вызову метода.

  3. Signal Event. Выбор нашего диспетчера, наиболее универсальный.

  4. Message Event. Resilience, scalability, availability, 10000+ RPS, high-load, k8s, «ставкинаспорт!» — если вы в таком контексте — это ваш выбор.

Способы делегирования
Способы делегирования

Подпроцесс. Шабашка для диспетчера

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

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

Call Activity. Экономия высокой ценой

Озвученная выше биллинговая метрика гласит: «a root process instance has no parent process instance». Экземпляры процессов, созданные через Call Activity (далее CA), в Camunda считаются дочерними и связаны с вызвавшим их родителем; таким образом, данный вариант маршрутизации дешевле — вызов обработчиков бесплатный. Более того, с точки зрения разработчика это проще, очевиднее, привычнее. Поэтому разберем все по порядку.

Семантика

Семантически CA говорит: «Родительскому процессу требуется твоя вычислительная логика, чтобы выполнить свою работу: возьми переданные параметры, останови его поток, верни ему результат». Ничего из этого роутеру не нужно, его работа заканчивается в момент отправки, результат обработки никак его не касается. Синхронный вызов для асинхронного сценария — ключевая ошибка подхода, все остальное — следствие.

Ответственность

Call Activities — handy shortcuts only within the boundary (Ruecker, Practical Process Automation)

Тезис про зону ответственности для CA применим в той же степени, что и для подпроцессов, но если в последнем случае разработчики это обычно даже на интуитивном уровне понимают, то CA часто используют, не задумываясь о границах контекста. Camunda в гайде по BPMN фактически приравнивает вызываемый через CA процесс к подпроцессу, сводя кейсы применения к modularization and reuse: Call Activity возникает не там, где два любых процесса надо связать, а там, где подпроцесс слишком разрастается или становится нужным где-то еще.

Масштабирование

Привычка использовать всюду синхронные вызовы в BPMN тянется со времен монолитов, которыми BPM-движки традиционно являлись. Zeebe, новое ядро Camunda 8, — Kubernetes-native распределенный кластерный движок с трехуровневым масштабированием.

Три уровня масштабирования Zeebe
Три уровня масштабирования Zeebe

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

Обработка ошибок

Самый очевидный для BPMN-разработчика практический недостаток CA — это трассировка ошибок. Как в коде, где исключение переходит вверх по стеку вызовов, пока его не обработают, BPMN-ошибка поднимается по скоупам, пока ее не поймает подходящий error catch event; через границу Call Activity она может быть поймана уже в родителе. Но если обработчик так и не встретился, инцидент оседает на бросившем элементе — внутри дочернего экземпляра. В дереве Operate он при этом подсвечивает и родителя вплоть до корня.

Просочившийся инцидент
Просочившийся инцидент

Чиновник в казначействе наспех состряпал неуместный ответ — а взыскали, помилуйте, с диспетчера: недельного жалованья, ни много ни мало, лишили, а тот и знать не знал, что там за письмо-то жалобное было. И не возразил ведь, смолчал.

Наблюдаемость и раздутый рантайм

Очевидное следствие синхронности CA — роутер живет столько же, сколько обработчик. Обработчики, напомню, выполняются в масштабе дней! В панели мониторинга процессов админ видит что-то вроде:

  • Активных инцидентов — 1594

  • Активных вопросов по бронированию — 933

  • Активных роутеров — 2675

Последний факт ничего не говорит о здоровье роутера: ключевая метрика процесса — в работе / завершен / ошибка — к нему не применима. Вместо анализа однозначных светофоров на дашборде приходится держать в голове особенности поведения каждого процесса.

Но висящий в работе экземпляр — это куда серьезнее, чем неэффективные метрики. Большинство BPM-движков строго разделяют активные инстансы (runtime) и завершенные (history). В Camunda 7 и Flowable это разные группы таблиц в реляционной БД (ACT_RU_* и ACT_HI_*); в Camunda 8 пошли дальше: рантайм держится в RocksDB на лидере партиции, а история сбрасывается во вторичное хранилище (Elasticsearch/OpenSearch, а с 8.9 еще и PostgreSQL). Захламление рантайма может ощутимо ударить по производительности и повысить стоимость репликации.

Зомби-апокалипсис

Если нас читает не разработчик, не архитектор, не админ, а, скажем, менеджер проекта — он, вероятно, пролистал по диагонали все технические эскапады, но совершенно точно запомнил две вещи:

  1. Синглтон роутер позволяет платить один раз за все выполненные за день маршрутизации.

  2. CA позволяет не платить за экземпляры вызываемых обработчиков.

Вот он, легальный хак лицензионного соглашения! Вот где заканчивается радуга и спрятано золото лепреконов! Захлебываясь от гениальности собственной идеи и потирая руки в предвкушении жирной премии, он бежит к аналитикам и судорожно воспроизводит на доске портрет нашего долгоживущего роутера с пририсованными Call Activities.

Помните, что каждое сообщение Kafka создает параллельный поток в роутере? Поток, который завершается за доли секунды трансляцией сигнала. Теперь же каждый такой подпоток ждет дни и часы. Висящие подпотоки не дают экземпляру завершиться, подписка открыта — «остановленный» роутер продолжает вычитывать топик и плодить новые подпотоки.

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

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

Сигнал. Предательство вендора

Защищать принятое решение всегда проще, когда уже рассказал о минусах соперника, поэтому кратко по всему озвученному.

  • Привычность: нет, для многих разработчиков BPMN асинхронная передача ощущается как потеря контроля

  • Семантика: да, fire-and-forget идеально описывает процедуру роутера

  • Ответственность: да, соблюдается loose coupling

  • Изолированность: полная — ошибки, состояния, рантайм находятся в своих исполняемых компонентах

Когда дело дошло до масштабирования, ответа в официальной документации я не нашел. На самом деле даже приведенное выше утверждение про единую для обоих экземпляров партицию при использовании Call Activity я вывел как следствие из общих рекомендаций в блоге Camunda.

Поиск информации о балансировке сигналов потребовал мини-расследования, в ходе которого я узнал про баг, который приводил к запуску процессов на КАЖДОЙ партиции. Регрессия просочилась в патч-релиз и в следующий минор, только представь себе, насколько нюансы реализации движка могут повлиять на то, как исполняется твой процесс.

Баг был устранен, конечно, довольно быстро, и Camunda радостно сообщила, что «сигнал теперь создает экземпляр только на одной партиции». Но на какой? Все еще не было ответа.

RTFM-метод провалился, поэтому, засучив рукава, поднимаем кластер и исследуем сами… Увы и ах, в этом раунде с CA ничья — несмотря на то, что сигнал транслируется на все партиции, создать экземпляр он может только там, откуда отправлен. Победа в общем зачете, а ограничение «все обработчики на одной партиции» для нашего сценария пока не страшно, ведь заявленная выше нагрузка до 10 RPS с головой выдерживается без шардирования.

Сообщения. Пресс-папье из чистого золота

В каморку к диспетчеру наведался давний однокашник, дослужившийся до главного координатора путей. Горделиво подмигнув, он ловким замахом достал из внутреннего кармана до блеска отполированный брусок и громким ударом приземлил его на стол перед обомлевшим работягой. «Это ж таким можно в день тыщи писем штамповать!» — «Да где ж у тебя их тыща, любезный? Ты его вон подыми попробуй — надорвешься!» — выхватил пресс-папье у потянувшегося к нему было диспетчера и, развернувшись, хихикая, вышел.

Message Event в Camunda — полноценная event-driven коммуникация со всеми преимуществами сигнала плюс масштабированием сверху. Балансировка сообщений осуществляется по хэшу ключа корреляции. Плата за это — потребность во внешнем транспорте.

Транспортировка сообщений Zeebe
Транспортировка сообщений Zeebe

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

Camunda балансирует на уровне Gateway: чтобы распределять нагрузку — необходимо выйти из брокера наружу.

Нашему роутеру, как выше говорилось, вполне хватает сигналов для вызова обработчиков, но в нем есть и приемщик сообщений: «Kafka-коннектор» — это не что иное, как message start event, настроенное на работу со вспомогательным компонентом кластера Camunda — Connector Runtime, в котором уже есть готовые Kafka-консьюмер и транслятор сообщений в Zeebe.

Раф остыл

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

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


1: термин «диспетчер» в этой статье выбран из художественных побуждений; автору известно, что описанный механизм не реализует паттерн «Dispatcher».

2: исходники доступны на GitHub.

3: FEEL — язык выражений из стандарта DMN, используемый в Camunda.

4: начиная с версии 8.9, в котором появился **Business ID**, это утверждение не актуально.