Zeebe
Zeebe is the process automation engine powering Camunda 8. While written in Java, you do not need to be a Java developer to use Zeebe.
With Zeebe you can:
- Define processes graphically in BPMN 2.0.
- Choose any gRPC-supported programming language to implement your workers.
- Build processes that react to events from Apache Kafka and other messaging platforms.
- Use as part of a software as a service (SaaS) offering with Camunda 8 or deploy with Docker and Kubernetes (in the cloud or on-premise) with Camunda 8 Self-Managed.
- Scale horizontally to handle very high throughput.
- Rely on fault tolerance and high availability for your processes.
- Export processes data for monitoring and analysis (currently only available through the Elasticsearch exporter added in Camunda 8 Self-Managed).
- Engage with an active community.
For documentation on deploying Zeebe as part of Camunda 8 Self-Managed, refer to the deployment guide.
Enterprise support for Zeebe
Paid support for Zeebe is available via either Camunda 8 Starter or Camunda 8 Enterprise plans. Customers can choose either plan based on their process automation requirements. Camunda 8 Enterprise customers also have the option of on-premise or private cloud deployment.
Additionally, regardless of how you are working with Zeebe and Camunda 8, you can always find support through the community.
Next steps
- Get familiar with technical concepts.
H Camunda представляет Open Source проект Zeebe для oркестрирования микросервисов инструментами BPMN в черновиках

Как и в оркестре, партитура имеет значение только для музыкантов и дирижёра. Для слушателей важно только само удовольствие от звучания.

Разделение монолитов на более мелкие части ведет к усложнению выполнения, мониторинга и конфигурации критических для бизнеса транзакций, которые теперь распределены между множеством микросервисов. Задача Zeebe является возвращение контроля над транзакциями в среде микросервисов.
Zeebe использует client/server aрхитектуру. Oсновными компонентами является брокер и Zeebe-клиент. Более подробное описание можно прочитать в документации по фраймворку.
Выступает DMN, дирижирует ZeeBe: как использовать бизнес-правила в микросервисах
Меня зовут Николай Первухин, я Senior Java Developer в Райффайзенбанке. Так сложилось, что, единожды попробовав бизнес-процессы на Camunda, я стал адептом этой технологии и стараюсь ее применять в проектах со сложной логикой. Действительно сама идея подкупает: рисуешь процесс в удобном GUI-редакторе (моделлере), а фреймворк выполняет эти действия по порядку, соблюдая большой спектр элементов нотации BPMN.
К тому же в Camunda есть встроенная поддержка еще одной нотации — DMN (Decision Model and Notation): она позволяет в простой и понятной форме создавать таблицы принятия решений по входящим наборам данных.

Но чего-то все же не хватает. Может, добавим немного скорости?
Почему ускоряем процессы
В банковской сфере бизнес-процессы широко используются там, где довольно часто встречаются длинные, с нетривиальной логикой процессы взаимодействия как с клиентом, так и между банковскими подсистемами.
Что обычно характеризует такие процессы:
- от момента создания до завершения процесса может пройти несколько дней;
- участвует большое количество сотрудников;
- осуществляется интеграция со множеством банковских подсистем.
Фреймворк Camunda прекрасно справляется с такими задачами, однако, заглянем под капот: там находится классическая база данных для осуществления транзакционности.
Но что если к логической сложности добавляется еще и требование к быстродействию? В таких случаях база данных рано или поздно становится «узким горлышком»: большое количество процессов начинает создавать блокировки, и это в конечном итоге приводит к замедлению.
Отличные новости: воспользуемся ZeeBe
В июле 2019 года было официально объявлено, что после двух лет разработки фреймворк ZeeBe готов к использованию на боевой среде. ZeeBe специально разрабатывался под задачи highload и, по утверждению автора, был протестирован при 10 000 процессов в секунду. В отличие от Camunda, ядро фреймворка ZeeBe принципиально не использует базу данных — из него убраны все вспомогательные подсистемы, в том числе и процессор правил DMN.
В случаях, когда DMN все же необходим, он может быть добавлен как отдельное приложение. Именно такую архитектуру мы и рассматриваем в данной статье.
- микросервис, инициирующий событие и запускающий процесс (event-handler);
- микросервис обработки бизнес-правил (rules-engine);
- микросервис, эмулирующий действия (action).
Данные микросервисы могут быть запущены в неограниченном количестве экземпляров для того, чтобы справиться с динамической нагрузкой.
Из оркестрации у нас:
- микросервис с брокером сообщений ZeeBe (zeebe);
- микросервис визуализации работающих процессов simplemonitor (zeebe-simple-monitor).
А присматривать за всеми микросервисами будет кластер k8s.
Схема взаимодействия

С точки зрения бизнес-логики в примере будет рассмотрен следующий бизнес-сценарий:
- из внешней системы происходит запрос в виде rest-обращения с передачей параметров;
- запускается бизнес-процесс, который «пропускает» входящие параметры через бизнес-правило;
- в зависимости от полученного решения из бизнес-правил запускается микросервис действий с различными параметрами.
Теперь поговорим подробнее о каждом микросервисе.
Микросервис zeebe
Данный микросервис состоит из брокера сообщений ZeeBe и экспортера сообщений для отображения в simple-monitor. Для ZeeBe используется готовая сборка, которую можно скачать с github. Подробно о сборке контейнера можно посмотреть в исходном коде в файле build.sh
Принцип ZeeBe — минимальное число компонентов, входящих в ядро, поэтому по умолчанию ZeeBe — это брокер сообщений, работающий по схемам BPMN. Дополнительные модули подключаются отдельно: например, для отображения процессов в GUI понадобится экспортер (доступны разные экспортеры, к примеру, в ElasticSearch, в базу данных и т.п.).
В данном примере возьмем экспортер в Hazelcast. И подключим его:
- добавим zeebe-hazelcast-exporter-0.10.0-jar-with-dependencies.jar в папку exporters;
- добавим в файл config/application.yaml следующие настройки:
exporters: hazelcast: className: io.zeebe.hazelcast.exporter.HazelcastExporter jarPath: exporters/zeebe-hazelcast-exporter-0.10.0-jar-with-dependencies.jar args: enabledValueTypes: "JOB,WORKFLOW_INSTANCE,DEPLOYMENT,INCIDENT,TIMER,VARIABLE,MESSAGE,MESSAGE_SUBSCRIPTION,MESSAGE_START_EVENT_SUBSCRIPTION" # Hazelcast port port: 5701
Данные активных процессов будут храниться в памяти, пока simplemonitor их не считает. Hazelcast будет доступен для подключения по порту 5701.
Микросервис zeebe-simplemonitor
Во фреймворке ZeeBe есть две версии GUI-интерфейса. Основная версия, Operate, обладает большим функционалом и удобным интерфейсом, однако, использование Operate ограничено специальной лицензией (доступна только версия для разработки, а лицензию для прода следует запрашивать у производителя).
Также есть облегченный вариант — simplemonitor (лицензируется по Apache License, Version 2.0)
Simplemonitor можно оформить тоже в виде микросервиса, который периодически подключается к порту hazelcast брокера ZeeBe и выгружает оттуда данные в свою базу данных.

Можно выбрать любую базу данных Spring Data JDBC, в данном примере используется файловая h2, где настройки, как и в любом Spring Boot приложении, вынесены в application.yml
spring: datasource: url: jdbc:h2:file:/opt/simple-monitor/data/simple-monitor-db;DB_CLOSE_DELAY=-1
Микросервис event-handler
Это первый сервис в цепочке, он принимает данные по rest и запускает процесс. При старте сервис осуществляет поиск файлов bpmn в папке ресурсов:
private void deploy() throws IOException < Arrays.stream(resourceResolver.getResources("classpath:workflow/*.bpmn")) .forEach(resource ->< try < zeebeClient.newDeployCommand().addResourceStream(resource.getInputStream(), resource.getFilename()) .send().join(); logger.info("Deployed: <>", resource.getFilename()); > catch (IOException e) < logger.error(e.getMessage(), e); >>); >
Микросервис имеет endpoint, и для простоты принимает вызовы по rest. В нашем примере передаются 2 параметра, сумма и лимит:
http://адрес-сервиса:порт/start?sum=100&limit=500
@GetMapping public String getLoad(@RequestParam Integer sum, @RequestParam Double limit) throws JsonProcessingException < Mapvariables = new HashMap<>(); variables.put("sum", sum); variables.put("limit", limit); zeebeService.startProcess(processName, variables); return "Process started"; >
Следующий код отвечает за запуск процесса:
public void startProcess(String processName, Map variables) throws JsonProcessingException
Сам процесс нарисован в специальной программе ZeeBe modeler (почти копия редактора Camunda modeler) и сохраняется в формате bpmn в папке workflow в ресурсах микросервиса. Графически процесс выглядит как:

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

Каждый дополнительный микросервис в данном примере будет использовать свой тип задач. Тип задач очень похож на очередь в Kafka: при возникновении задач к нему могут подключаться подписчики — worker’ы.
После старта процесс продвинется на один шаг, и появится сообщение типа DMN.
Микросервис rules-engine
Благодаря прекрасной модульной архитектуре Camunda есть возможность использовать в своем приложении (отдельно от самого фреймворка Camunda) движок правил принятия решения.
Для его интеграции с вашим приложением достаточно добавить его в зависимости maven:
org.camunda.bpm.dmn camunda-engine-dmn $
Сами правила создаются в специальном графическом редакторе Camunda modeler. Одна диаграмма DMN имеет два вида отображения.
Entity Relation Diagram (вид сверху) показывает зависимости правил друг от друга и от внешних параметров:

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

Как видно из примера выше, бизнес-правило представляется в виде таблицы, в которой перечислены входящие и результирующие параметры. Инструмент достаточно богатый, и можно использовать различные методы сравнения, диапазоны, множества, несколько типов политик правил (первое совпадение, множественное, последовательность по диаграмме и т.п.). Такая диаграмма сохраняется в виде файла dmn.
Посмотрим на примере, где такой файл располагается в папке dmn-models в ресурсах микросервиса. Для регистрации диаграммы при старте микросервиса происходит его однократная загрузка:
public void init() throws IOException < Arrays.stream(resourceResolver.getResources("classpath:dmn-models/*.dmn")) .forEach(resource ->< try < logger.debug("loading model: <>", resource.getFilename()); final DmnModelInstance dmnModel = Dmn.readModelFromStream((InputStream) Resources .getResource("dmn-models/" + resource.getFilename()).getContent()); dmnEngine.parseDecisions(dmnModel).forEach(decision -> < logger.debug("Found decision with id '<>' in file: <>", decision.getKey(), resource.getFilename()); registry.put(decision.getKey(), decision); >); > catch (IOException e) < logger.error("Error parsing dmn: <>", resource, e); > >); >
Для того, чтобы подписаться на сообщения от ZeeBe, требуется осуществить регистрацию worker’а:
private void subscribeToDMNJob() < zeebeClient.newWorker().jobType(String.valueOf(jobWorker)).handler( (jobClient, activatedJob) -> < logger.debug("processing DMN"); final String decisionId = readHeader(activatedJob, DECISION_ID_HEADER); final Mapvariables = activatedJob.getVariablesAsMap(); DmnDecisionResult decisionResult = camundaService.evaluateDecision(decisionId, variables); if (decisionResult.size() == 1) < if (decisionResult.get(0).containsKey(RESULT_DECISION_FIELD)) < variables.put(RESULT_DECISION_FIELD, decisionResult.get(0).get(RESULT_DECISION_FIELD)); >> else < throw new DecisionException("Нет результата решения."); >jobClient.newCompleteCommand(activatedJob.getKey()) .variables(variables) .send() .join(); > ).open(); >
В данном коде осуществляется подписка на событие DMN, вызов модели правил при получении сообщения от ZeeBe и результат выполнения правила сохранятся обратно в бизнес-процесс в виде переменной result (константа RESULT_DECISION_FIELD ).
Когда данный микросервис отчитывается ZeeBe о выполнении операции, бизнес-процесс переходит к следующему шагу, где происходит выбор пути в зависимости от выполнения условия, заданного в свойствах стрелочки:

Микросервис action
Микросервис action совсем простой. Он также осуществляет подписку на сообщения от ZeeBe, но другого типа — action .

В зависимости от полученного результата будет вызван один и тот же микросервис action, но с различными параметрами. Данные параметры задаются в закладке headers:

Также передачу параметров можно сделать и через закладку Input/Output, тогда параметры придут вместе с переменными процесса, но передача через headers является более «каноничной».
Посмотрим на получение сообщения в микросервисе:
private void subscribe() < zeebeClient.newWorker().jobType(String.valueOf(jobWorker)).handler( (jobClient, activatedJob) -> < logger.debug("Received message from Workflow"); actionService.testAction( activatedJob.getCustomHeaders().get(STATUS_TYPE_FIELD), activatedJob.getVariablesAsMap()); jobClient.newCompleteCommand(activatedJob.getKey()) .send() .join(); >).open(); >
Здесь происходит логирование всех переменных бизнес-процесса:
public void testAction(String statusType, Map variables) < logger.info("Event Logged with statusType <>", statusType); variables.entrySet().forEach(item -> logger.info("Variable <> = <>", item.getKey(), item.getValue())); >
Исходный код
Весь исходный код прототипа можно найти в открытом репозитории GitLab.
Компиляция образов Docker
Все микросервисы проекта собираются командой ./build.sh
Для каждого микросервиса есть различный набор действий, направленных на подготовку образов docker и загрузки этих образов в открытые репозитории hub.docker.com.
Загрузка микросервисов в кластер k8s
Для развертывания в кластере необходимо сделать следующие действия:
- Создать namespace в кластере kubectl create namespace zeebe-dmn-example
- Создать config-map общих настроек
kind: ConfigMap apiVersion: v1 metadata: name: shared-settings namespace: zeebe-dmn-example data: shared_servers_zeebe:
Далее создаем два персистентных хранилища для хранения данных zeebe и simplemonitor . Это позволит осуществлять перезапуск соответствующих подов без потери информации:
kubectl apply -f zeebe—sm-volume.yml
kubectl apply -f zeebe-volume.yml
Yml-файлы находятся в соответствующих проектах:

Теперь осталось последовательно создать поды и сервисы. Указанные yml-файлы находятся в корне соответствующих проектов.
kubctl apply -f zeebe-deployment.yml
kubctl apply -f zeebe-sm-deployment.yml
kubctl apply -f event-handler-deployment.yml
kubctl apply -f rules-engine-deployment.yml
kubctl apply -f action-deployment.yml
Смотрим, как отображаются наборы подов в кластере:

И мы готовы к тестовому запуску!
Запуск тестового процесса
Запуск процесса осуществляется открытием в браузере соответствующий URL. К примеру, сервис event-handler имеет сервис с внешним IP и портом 81 для быстрого доступа.
Далее можно проверить отображение процесса в simplemonitor . У данного микросевиса тоже есть внешний сервис с портом 82.

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

Поделюсь, какими ресурсами пользовался для подготовки прототипа
- https://www.youtube.com/watch?v=Q3tLKV-6J3c
- https://github.com/zeebe-io/zeebe-dmn-worker
- https://github.com/berndruecker/zeebe-camunda-dmn/blob/master/README.md
- https://zeebe.io/blog/2019/08/zeebe-horizontal-scalability/
Небольшое послесловие вместо итогов
Из плюсов:
- архитектура разработанного прототипа получилась гибкой и расширяемой. Можно добавлять любое количество микросервисов для обработки сложной логики;
- простая нотация BPMN и DMN позволяет привлекать аналитиков и бизнес к обсуждению сложной логики;
- Zeebe показал себя как очень быстрый фреймворк, задержки на получение/отправку сообщений практически отсутствуют;
- Zeebe изначально разрабатывался для работы в кластере, в случае необходимости можно быстро нарастить мощности;
- без ZeeBe Operate можно вполне обойтись, Simple-Monitor отвечает минимальным требованиям.
Из минусов:

- хотелось бы иметь возможность редактирования DMN непосредственно в ZeeBe modeler (как это реализовано в Camunda), на данный момент, приходится использовать оба моделлера;
- к сожалению, только в Enterprise версии Camunda есть возможность просмотра пути, по которому принималось решение:
Это очень полезная функция при отладке правил. В Community версии приходится добавлять дополнительное поле типа output для логирования, либо разработать свое решение визуализации.
При развертывании прототипа в реальных условиях в кластере k8s необходимо добавить ограничения по ресурсам (CPU и RAM), также нужно будет подобрать лучшую систему хранения исторических данных.
Где применять такие технологии:
- как оркестрация внутри одной команды или продукта в виде перекладывания сложной логики на диаграммы BPMN / DMN;
- в сфере, где идет потоковая обработка данных с большим количеством интеграций. В банке это может быть проведение или проверка транзакций, обработка данных из внешних систем или просто многоэтапная трансформация данных;
- как частичная альтернатива существующего стека ESB или Kafka для интеграции между командами.
Коллеги, понимаю, что есть множество разных технологий и подходов. Буду рад, если поделитесь в комментариях вашим опытом: как вы решаете подобные задачи?
Zeebe и Camunda: сравниваем известные BPM-системы под высокими нагрузками 14.12.2021 13:31
Всем привет! Меня зовут Николай Первухин, я Senior Java Developer в Райффайзенбанке. В последнее время я активно занимаюсь BPM-системами Camunda и Zeebe (основа Camunda-cloud). Если вы, как и я, с ходу не можете ответить на вопрос, кто быстрее — Camunda или Zeebe, насколько, и в каких случаях они могут тормозить, — то вам будет очень интересно прочитать эту статью.
В этом материале я попытаюсь оценить производительность систем Camunda и Zeebe с различными параметрами, коснусь классических систем по workflow — Apache Camel и Spring Integration, а также постараюсь предложить гибридное решение для повышения производительности.
Часто со стороны бизнеса можно услышать вопрос, что будет, если поток данных возрастет — справится ли наша система с этим. Вот и мне стало интересно, так что давайте немного поэкспериментируем на простом workflow: сделаем несколько действий и осуществим rest-вызов в каждом из них.
Вызывать будем статический тестовый файл test.json с содержимым: <> , который будет выдавать локальный nginx.
Немного об оборудовании: в моем распоряжении 4х ядерный Intel ® Core ™ i7–4770K CPU @ 3.50GHz, SSD диск и 24Gb памяти.
Итак, наши кандидаты:
Apache-camel — чемпион, проверенный временем
Признаюсь, это вообще не BPM. Но этот фреймворк настолько популярен, что стоит начать рассказ для человека, который не в теме BPMN, и его глаза вспыхивают пониманием — это же про Apache Camel! Действительно, задачу он решает похожую, поэтому давайте рассмотрим его с должным уважением.
Вот воркфлоу, который будем исполнять:
camelContext.addRoutes(new RouteBuilder() < @Override public void configure() throws Exception < from("direct:tuneBPMN") .bean("restExecutorService", "doRestRequest") .bean("restExecutorService", "doRestRequest"); >>);
Будем производить тестовый REST-вызов:
public void doRestRequest() < final String result = restTemplate.getForObject(restRequestUrl + "?v java">
Ссылка проекта на gitlab
Camunda — надежный, как DasAuto
Пока рассмотрим самый быстрый вариант Camunda — мы отключим историю вообще, и будем использовать in-memory h2 базу данных.
Наш «сложный» процесс выглядит так:
Код для RestDelegate достаточно прост:
@Component public class RestDelegate implements JavaDelegate < private static final Logger logger = LoggerFactory.getLogger(RestDelegate.class); private static final Random random = new Random(); . @Override public void execute(DelegateExecution execution) throws Exception < final String result = restTemplate.getForObject(restRequestUrl + "?v https://habr.com/img/image-loader.svg" alt="image-loader.svg" /> Судя по графику, при одном инcтансе результаты должны быть в пределе 1 тыс. процессов в секунду.
Будем использовать данный процесс:
Наш воркер:
@ZeebeWorker(type = "restJob") public void handleJob(JobClient jobClient, ActivatedJob activatedJob) < final String result = restTemplate.getForObject(restRequestUrl + "?v pgsql">docker run -d --name zeebe -e ZEEBE_BROKER_CLUSTER_PARTITIONSCOUNT=24 -p 26500:26500 camunda/zeebe
Если вам потребуется вынести папку с данными, используйте параметр:
-v ваша_локальная_папка_с_данными:/usr/local/zeebe/data
Небольшой тюнинг (в рамках 1 иснтанса):
ZEEBE_BROKER_CLUSTER_PARTITIONSCOUNT — число партиций кластера (по умолчанию это 1) определяет, на какое количество частей или папок разбивать rockdb для дальнейшей синхронизации. В зависимости от скорости работы SSD или жесткого диска, можно подобрать лучший параметр. В данном случае устанавливается 24, потому что дальнейшее увеличение дает скорее ухудшение.
При этом два дополнительных параметра не оказывают большую роль, так как процесс относительно простой:
Мы ставим их по умолчанию. Клиент для Zeebe с воркером можно найти на gitlab.
Даем нагрузку
Для замера нагрузки будет произведен запуск 100 тыс. процессов, при этом с интервалом 10 тыс. мы планируем замерять время выполнения. Прошу простить меня за неточность, но начало отсчета будем фиксировать не с 0, а с 1 тыс. процессов. Это важно, чтобы оценить, насколько быстро происходит прогрев системы. Результат нагрузки:
Объясню полученные результаты:
- Apache Camel и Spring Integration вне конкуренции, когда сложность процесса небольшая и нет изменчивости. Так что +100 очков команде Apache Camel, они даже немного обошли Spring Integration. Здесь было важно показать, что условные 4,2 тыс. процессов — наилучший возможный результат на 1 инстансе.
- Camunda без истории и in-memory показывает достаточно хорошую производительность.
- Zeebe даже без экспорта на одном инстансе более чем скромен — как и ожидалось.
Все вышеописанные результаты получены при полном отсутствии логирования — это как раз тот случай, когда нам не требуется контролировать результаты.
Логирование, экспортирование
Теперь нам требуется посмотреть, как проходил процесс, какие таски запускались и какие были переменные. В этом случае мы должны включить логирование процесса — или экспортирование, говоря в терминах Zeebe.
Camunda
Camunda поддерживает транзакционные SQL-базы. Канонической для Camunda базой является PostgreSQL — это наиболее частый случай использования, которую будем использовать в нашем эксперименте.
Базу данных подключаем в виде готового докер контейнера, тут все стандартно:
docker run -d --name postgres -e POSTGRES_PASSWORD=postgres -p 5432:5432 postgres
Если требуется вынести папку с данными, используйте параметр:
-v ваша_локальная_папка_с_данными:/var/lib/postgresql/data
Дальше мы переключаем профиль в pure-camunda на PostgreSQL.
Zeebe
В Zeebe подключаем готовый экспортер для ElasticSearch. Данные из Elastic потом можно увидеть в Zeebe: Operate (аналог Camunda Cockpit). Подключаем ElasticSearch тоже через docker-контейнер. Для простоты мы будем использовать сеть рабочей машины, чтобы Zeebe и ElasticSearch могли общаться напрямую.
docker run -d --name elastic --network host -e "discovery.type=single-node" elasticsearch:7.14.2
Если требуется вынести папку с данными, используйте параметр:
-v ваша_локальная_папка_с_данными:/usr/share/elasticsearch/data
В самом Zeebe подключаем стандартный экспортер в ElasticSearch:
docker run -d --name zeebe --network host -e ZEEBE_BROKER_EXPORTERS_ELASTICSEARCH_CLASSNAME=io.camunda.zeebe.exporter.ElasticsearchExporter -e ZEEBE_BROKER_EXPORTERS_ELASTICSEARCH_ARGS_URL=http://localhost:9200 -e ZEEBE_BROKER_EXPORTERS_ELASTICSEARCH_ARGS_BULK_SIZE=1000 -e ZEEBE_BROKER_CLUSTER_PARTITIONSCOUNT=24 camunda/zeebe
Из новых параметров среды:
- ZEEBE_BROKER_EXPORTERS_ELASTICSEARCH_CLASSNAME — тип экспортера
- ZEEBE_BROKER_EXPORTERS_ELASTICSEARCH_ARGS_URL — адрес ElasticSearch
- ZEEBE_BROKER_EXPORTERS_ELASTICSEARCH_ARGS_BULK_SIZE — размер пачки экспортируемых сообщений
А вот и результаты:
Объясню полученные результаты:
- Результат с Camunda изменился сильно, ухудшение почти в три раза. Более того, идет деградирование с увеличением данных в базе. Это связано, в основном, с ребилдом индексов, которые Camunda активно использует.
- А вот у Zeebe значения почти не изменились — экспорт оказал минимальное влияние на время исполнения. Результат скромный, но стабильный.
Экспортируемые данные накапливаются и отправляются в систему хранения пачками по таймеру, либо по мере накопления такой пачки. Здесь достигается двойной эффект:
- Экспорт вынесен отдельно от процессинга. Если данные вдруг перестанут экспортироваться — процессинг все равно идет своим чередом.
- Добавление данных большими объемами (batch mode) почти в любых системах хранения эффективнее, чем по одной записи.
Создание гибрида bpm-light
Давайте подумаем, как создать гибридную систему, чтобы она была:
- Такая же масштабируемая, как Zeebe
- Быстрая как in-memory Camunda. Лучше быстрее, чем Apache Camel, но это мечты 🙂
- С возможностью экспорта данных, чтобы визуально было видно, как отработал процесс.
Помните, как в фильме «Марсианин» снимали с ракеты все ненужное?
Избавляемся от тяжелых компонентов, оставляем только движок Camunda.
В состав pom включаем компоненты:
- Spring boot, spring data и web-services
- Только модуль движка Camunda:
docker run -d --name zeebe --network host -e ZEEBE_BROKER_EXPORTERS_ELASTICSEARCH_CLASSNAME=io.camunda.zeebe.exporter.ElasticsearchExporter -e ZEEBE_BROKER_EXPORTERS_ELASTICSEARCH_ARGS_URL=http://localhost:9200 -e ZEEBE_BROKER_EXPORTERS_ELASTICSEARCH_ARGS_BULK_SIZE=1000 -e ZEEBE_BROKER_CLUSTER_PARTITIONSCOUNT=24 camunda/zeebe
- Драйверы баз данных h2, clickhouse
- Средство развертывания базы flyway
Из интересного (надеюсь):
Конфигурация движка Camunda. Используем StandaloneInMemProcessEngineConfiguration — думаю, название говорит само за себя:
@Bean public ProcessEngine processEngine(@Autowired RestDelegate restDelegate, @Autowired ExporterRegistryService exporterRegistryService) < final StandaloneInMemProcessEngineConfiguration config = new StandaloneInMemProcessEngineConfiguration(); // используем полное логирование процесса, full — включая переменные config.setHistory("full"); config.setDatabaseSchemaUpdate("create"); // отключаем сбор метрик config.setMetricsEnabled(false); config.setDbMetricsReporterActivate(false); // Добавление бина делегата config.setBeans(new HashMap<>()); config.getBeans().put("restDelegate", restDelegate); // Обработчик событий процесса (старт процесса, этапы процесса, конец процесса, изменение переменных и тп) config.setHistoryEventHandler(new CustomHistoryHandler(exporterRegistryService)); config.setIdGenerator(new StrongUuidGenerator()); return config.buildProcessEngine(); >
Разрабатываем свой собственный регистратор исторических событий CustomHistoryHandler :
@Override public void handleEvent(HistoryEvent historyEvent) < if (historyEvent.getClass().equals(HistoricProcessInstanceEventEntity.class)) < final HistoricProcessInstanceEventEntity event = (HistoricProcessInstanceEventEntity) historyEvent; // логируем начало и окончание процесса if ("start".equals(event.getEventType())) < exporterRegistryService.registerStartProcess(event.getStartTime(), event.getProcessDefinitionKey(), event.getProcessDefinitionId(), event.getBusinessKey(), event.getProcessInstanceId(), event.getSuperProcessInstanceId()); return; >if ("end".equals(event.getEventType())) < exporterRegistryService.registerEndProcess(event.getEndTime(), event.getProcessDefinitionKey(), event.getProcessDefinitionId(), event.getBusinessKey(), event.getProcessInstanceId(), event.getSuperProcessInstanceId()); return; >>
Аналогично перехватываем начало и завершение задач в процессе, а также создание и изменение переменных. Здесь очень важным моментом является, что метод полностью подменяет сохранение исторических данных — это значит, что в памяти исторические данные мы больше не храним.
Идею с экспортером данных честно позаимствуем у Zeebe. Будем накапливать события в LinkedBlockingQueue .
// Накапливаем события в LinkedBlockingQueue по каждому типу событий public void export(ExportRecordType type, ExportRecord exportRecord) < if (enabled) < if (! this.exportLog.containsKey(type)) < this.exportLog.put(type, new LinkedBlockingQueue<>()); > this.exportLog.get(type).add(exportRecord); // в случае переполнения пачки, мы запускаем запись в систему хранения if (this.exportLog.get(type).size() >= this.batchSize) < flush(type); >> > // Полная запись всех накопленных событий в систему хранения private void flushBatchTask() < fullFlush(); scheduleFlushBatchTask(); >private void fullFlush() < exportLog.entrySet().stream() .forEach( x -> < if (x.getValue() != null && ! x.getValue().isEmpty()) < flush(x.getKey()); >>); > // Создаем отложенное событие запуска по таймеру private void scheduleFlushBatchTask() < executor.schedule(this::flushBatchTask, interval, TimeUnit.SECONDS); >// Запускаем непосредственный слив данных в систему хранения private void flush(ExportRecordType type) < if (! isRunning) < isRunning = true; exporter.export(type, exportLog.get(type)); isRunning = false; >else < logger.debug("Слишком большая нагрузка, экспорт не успевает"); >>
Метод *.export вынесен как интерфейс для поддержки различных систем хранения. Так было сделано в Zeebe, мы же просто позаимствуем этот принцип. Переключать систему можно профилем spring-boot.
Для эксперимента я сделал три различных экспортера:
- В каноническую базу данных H2. На ней легко тестировать, не требует установки
- В ElasticSearch. Наверное, это один из самых быстрых индексаторов json-данных, который можно рассматривать как быстрое и масштабируемое хранилище
- В Clickhouse. Эта база данных прекрасно себя зарекомендовала при быстрой вставке больших массивов данных, по описанию — то, что нам нужно.
Рассмотрим на примере h2:
@Override public void export(ExportRecordType exportRecordType, BlockingQueue records) < .. // определяем команду sql для вставки или изменений данных и набор параметров switch (exportRecordType) < .. case PROCESS: queryString = "insert into PROCESS " + "(CREATED, PROCESSDEFINITIONKEY, PROCESSDEFINITIONID, BUSINESSKEY, PROCESSINSTANCEID, " + "SUPERPROCESSINSTANCEID, LIFECYCLETYPE,ENDDATE) values "; valuesString = "(. )"; break; .. >// вычитываем данные из LinkedBlockingQueue и создаем один большой запрос с параметрами while (! records.isEmpty()) < final ExportRecord recordElement = records.poll(); itemsAdded++; switch (exportRecordType) < .. case PROCESS: Process process = (Process) recordElement; parameters.add(process.getDate()); parameters.add(process.getProcessDefinitionKey()); parameters.add(process.getProcessDefinitionId()); parameters.add(process.getBusinessKey()); parameters.add(process.getProcessInstanceId()); parameters.add(process.getSuperProcessInstanceId()); parameters.add(process.getLifecycleType().name()); parameters.add(process.getEndDate()); break; .. >> if (itemsAdded > 0 && ! parameters.isEmpty()) < String query = queryString + (valuesString + ",").repeat(itemsAdded); query = query.substring(0, query.length() - 1); execute(query, parameters); >>
Сам же метод execute добавляет параметры в jdbc-запрос в зависимости от типа и выполняет запрос:
private void execute(String query, List parameters) < try (PreparedStatement statement = databaseConnectionService.getConnection() .prepareStatement(query)) < for (int i=0; iif (parameters.get(i) instanceof Long) < statement.setLong(i+1, (Long) parameters.get(i)); >.. if (parameters.get(i) == null) < statement.setString(i+1, null); >> // выполняем запрос statement.execute(); > catch (Exception ex) < logger.error(ex.getMessage(), ex); >>
Фронт у системы тоже должен быть. Для этого создадим rest-api и различные методы извлечения данных для каждой системы хранения.
Фронт сделан минималистично, он позволяет:
- Просматривать и устанавливать новые bpmn-процессы
- Смотреть активные и завершенные инстансы
- Видеть путь прохождения инстанса, все переменные и тайминги по каждой из задач.
Проект фронта использует jsf primefaces и bpmn.io для рендеринга самих процессов.
Если интересно, то можно ознакомиться с исходным кодом здесь:
Код серверной части
Код фронт части
Вернемся к нагрузке и дадим такую же нагрузку на этот проект для каждого из методов хранения:
- На удивление прекрасно показал себя Elastic, неожиданно обогнав Clickhouse. Результаты h2 тоже впечатлили бы, если бы не одно большое но: без индексов на такой базе выборки начинают тупить уже после 5 тыс. процессов, да и фронт уже перестает отвечать. В общем, это хорошая база для тестов, но не для реальной жизни.
- Важной отличительной частью такого подхода является то, что такой гибрид можно запустить в большом количестве инстансов и поставить за распределителем нагрузки (eureka, nginx и тп). Получится, что вы сможете без каких-либо ограничений горизонтально масштабировать решение.
- Главный минус — неперсистентность активных инстансов. Сами процессы идут в оперативной памяти. В случае аварийной перезагрузки исторические данные по работающим экземплярам могут не успеть экспортироваться в хранилище данных.
Итоги
Итак, пришло время ответить на вопросы в начале статьи:
Кто быстрее — Camunda или Zeebe?
- При равных условиях выигрывает Camunda. Но при марафоне с нарастанием базы рано или поздно Camunda упирается в деградацию и начинает замедляться. Zeebe из коробки с дефолтными настройками показывает очень скромные результаты, зато практически не имеет пределов для масштабирования.
А насколько, и в каких случаях?
- При равных условиях в виде отсутствия экспорта и истории, Camunda почти в 15 раз быстрее, чем Zeebe. При подключении истории производительность Camunda резко падает — почти в три раза, а эффективный метод экспортирования данных у Zeebe такой просадки не дает.
- Для совсем небольших процессов без необходимости хранить историю шагов даже Camunda не является целесообразной. В этих случаях лучше выбрать Apache Camel или Spring Integration — такое решение получится как минимум в два раза быстрее.
- Zeebe нужна для совсем крупных задач и высоконагруженных процессов (коротких и нагруженных вычислениями), чтобы быть экономически целесообразным. Для обеспечения высокой производительности требуется кластер с большим количеством инстансов и быстрыми дисками.
- В случае, когда бюджет ограничен, но есть желание получить скорость и масштабируемость, лучше создать аналогичный гибрид — отучить Camunda от использования классической базы данных.
За счет чего может тормозить?
- Camunda тормозит из-за своей архитектуры — транзакционная запись каждого действия в базу данных. Это крутая фича, она позволяет при возникновении ошибок процессу откатываться на нужный этап, но это совсем не про hi-load.
- Zeebe тормозит из-за своего ядра — распределенной базы данных rockdb. Для нее нужны очень быстрые диски и большое количество партиций и инстансов.
Ссылка на код группы проектов
Коллеги, а как вы оптимизируете производительность Camunda и Zeebe? Поделитесь в комментариях!