Camunda Zeebe
Экспериментальный модуль
Экспериментальный модуль является полностью рабочим и протестированным, но требует дополнительной апробации и аналитики по использованию,
по этой причине его API может потенциально претерпеть незначительные изменения перед полной готовностью.
Модуль подключает клиент Camunda 8 (Zeebe) и создает
исполнителей заданий для внешнего оркестратора процессов. В Kora такой исполнитель объявляется обычным компонентом:
метод с аннотацией @JobWorker получает переменные процесса, выполняет работу и возвращает результат, который будет
передан обратно в Zeebe.
Подключение¶
Зависимость build.gradle:
Модуль:
Зависимость build.gradle.kts:
Модуль:
Конфигурация¶
Пример полной конфигурации клиента, описанной в классе ZeebeClientConfig (указаны примеры значений или значения по умолчанию):
zeebe {
client {
executionThreads = 2 //(1)!
keepAlive = "45s" //(2)!
tls = true //(3)!
certificatePath = "/file/path/to/cert.crt" //(4)!
initializationFailTimeout = "15s" //(5)!
grpc {
url = "grpc://localhost:8090" //(6)!
ttl = "1h" //(7)!
maxMessageSize = "4Mib" //(8)!
retryPolicy {
enabled = true //(9)!
attempts = 5 //(10)!
delay = "100ms" //(11)!
delayMax = "5s" //(12)!
step = 3.0 //(13)!
}
}
rest {
url = "http://localhost:8080" //(14)!
}
deployment {
resources = "classpath:bpm" //(15)!
timeout = "45s" //(16)!
}
telemetry {
logging {
enabled = false //(17)!
}
metrics {
enabled = true //(18)!
slo = [ 1, 10, 50, 100, 200, 500, 1000, 2000, 5000, 10000, 20000, 30000, 60000, 90000 ] //(19)!
tags = { // (20)!
"key1" = "value1"
"key2" = "value2"
}
}
tracing {
enabled = true //(21)!
attributes = { // (22)!
"key1" = "value1"
"key2" = "value2"
}
}
}
}
}
- Максимальное количество потоков для исполнителей заданий (по умолчанию: количество ядер процессора, но не меньше
2) - Время без активности чтения перед отправкой проверки
KeepAlive(по умолчанию:45s) - Использовать ли
TLSпри подключении (по умолчанию:true) - Файловый путь до сертификата для подключения; если не указан, используется системный сертификат (по умолчанию не указано, необязательно)
- Максимальное время ожидания проверки доступности топологии при запуске клиента (по умолчанию не указано, необязательно)
URLдля подключения поgRPC(обязательная, по умолчанию не указано)- Время, в течение которого сообщение должно храниться на брокере при отправке по
gRPC(по умолчанию:1h) - Максимальный размер входящего сообщения по
gRPC(по умолчанию:4Mib) - Включена ли политика повторов для
gRPC-соединения (по умолчанию:true) - Количество попыток (по умолчанию:
5) - Начальная задержка между попытками (по умолчанию:
100ms) - Максимальная задержка между попытками (по умолчанию:
5s) - Коэффициент увеличения задержки между попытками (по умолчанию:
3.0) URLдля подключения кREST-адресуZeebe; если указан, клиент предпочитаетRESTвместоgRPCдля поддерживаемых операций (обязательнаявнутри необязательной секцииrest, по умолчанию не указано)- Пути для поиска ресурсов, которые будут загружены в оркестратор после запуска (по умолчанию:
[]) - Максимальное время ожидания загрузки ресурсов (по умолчанию:
45s) - Включает логирование модуля (по умолчанию:
false) - Включает метрики модуля (по умолчанию:
true) - Настройка SLO для метрик (по умолчанию:
ru.tinkoff.kora.telemetry.common.TelemetryConfig.MetricsConfig#DEFAULT_SLO) - Настройка тегов для метрик (по умолчанию:
{}) - Включает трассировку модуля (по умолчанию:
true) - Настройка атрибутов для трассировки (по умолчанию:
{})
zeebe:
client:
executionThreads: 2 #(1)!
keepAlive: "45s" #(2)!
tls: true #(3)!
certificatePath: "/file/path/to/cert.crt" #(4)!
initializationFailTimeout: "15s" #(5)!
grpc:
url: "grpc://localhost:8090" #(6)!
ttl: "1h" #(7)!
maxMessageSize: "4Mib" #(8)!
retryPolicy:
enabled: true #(9)!
attempts: 5 #(10)!
delay: "100ms" #(11)!
delayMax: "5s" #(12)!
step: 3.0 #(13)!
rest:
url: "http://localhost:8080" #(14)!
deployment:
resources: "classpath:bpm" #(15)!
timeout: "45s" #(16)!
telemetry:
logging:
enabled: false #(17)!
metrics:
enabled: true #(18)!
slo: [ 1, 10, 50, 100, 200, 500, 1000, 2000, 5000, 10000, 20000, 30000, 60000, 90000 ] #(19)!
tags: #(20)!
key1: value1
key2: value2
tracing:
enabled: true #(21)!
attributes: #(22)!
key1: value1
key2: value2
- Максимальное количество потоков для исполнителей заданий (по умолчанию: количество ядер процессора, но не меньше
2) - Время без активности чтения перед отправкой проверки
KeepAlive(по умолчанию:45s) - Использовать ли
TLSпри подключении (по умолчанию:true) - Файловый путь до сертификата для подключения; если не указан, используется системный сертификат (по умолчанию не указано, необязательно)
- Максимальное время ожидания проверки доступности топологии при запуске клиента (по умолчанию не указано, необязательно)
URLдля подключения поgRPC(обязательная, по умолчанию не указано)- Время, в течение которого сообщение должно храниться на брокере при отправке по
gRPC(по умолчанию:1h) - Максимальный размер входящего сообщения по
gRPC(по умолчанию:4Mib) - Включена ли политика повторов для
gRPC-соединения (по умолчанию:true) - Количество попыток (по умолчанию:
5) - Начальная задержка между попытками (по умолчанию:
100ms) - Максимальная задержка между попытками (по умолчанию:
5s) - Коэффициент увеличения задержки между попытками (по умолчанию:
3.0) URLдля подключения кREST-адресуZeebe; если указан, клиент предпочитаетRESTвместоgRPCдля поддерживаемых операций (обязательнаявнутри необязательной секцииrest, по умолчанию не указано)- Пути для поиска ресурсов, которые будут загружены в оркестратор после запуска (по умолчанию:
[]) - Максимальное время ожидания загрузки ресурсов (по умолчанию:
45s) - Включает логирование модуля (по умолчанию:
false) - Включает метрики модуля (по умолчанию:
true) - Настройка SLO для метрик (по умолчанию:
ru.tinkoff.kora.telemetry.common.TelemetryConfig.MetricsConfig#DEFAULT_SLO) - Настройка тегов для метрик (по умолчанию:
{}) - Включает трассировку модуля (по умолчанию:
true) - Настройка атрибутов для трассировки (по умолчанию:
{})
Предоставляемые метрики модуля описаны в разделе Справочник метрик.
Развертывание ресурсов¶
Если в deployment.resources указаны пути, модуль во время запуска находит ресурсы в classpath и развертывает их в
Zeebe через компонент ZeebeResourceDeployment. Развертываются как BPMN-процессы, так и DMN-решения, найденные в
настроенных расположениях. Поддерживаются только пути с префиксом classpath:, например classpath:bpm; другие
расположения логируются и пропускаются.
Разместите развертываемые ресурсы в соответствующей директории classpath:
- Одно или несколько расположений в classpath для поиска
BPMN/DMN-ресурсов (одно значение или список)
Клиент¶
Модуль создает компонент ZeebeClient, который можно внедрять в собственные сервисы, если нужно вручную запускать
процессы, публиковать сообщения или выполнять другие команды Zeebe.
Например, чтобы запустить новый экземпляр процесса:
@Component
public final class ProcessStarter {
private final ZeebeClient client;
public ProcessStarter(ZeebeClient client) {
this.client = client;
}
public void start() {
ProcessInstanceEvent event = client.newCreateInstanceCommand()
.bpmnProcessId("demo") //(1)!
.latestVersion() //(2)!
.variables("{\"startId\":\"42\"}") //(3)!
.send()
.join(); //(4)!
}
}
- Идентификатор
BPMN-процесса, который нужно запустить - Запуск последней развернутой версии процесса
- Начальные переменные процесса в виде
JSON-строки (также принимаютсяMapили@Json-объект) - Отправить команду и заблокироваться до подтверждения от
Zeebe(используйте возвращаемыйCompletionStageдля неблокирующего вызова)
@Component
class ProcessStarter(private val client: ZeebeClient) {
fun start() {
val event = client.newCreateInstanceCommand()
.bpmnProcessId("demo") //(1)!
.latestVersion() //(2)!
.variables("""{"startId":"42"}""") //(3)!
.send()
.join() //(4)!
}
}
- Идентификатор
BPMN-процесса, который нужно запустить - Запуск последней развернутой версии процесса
- Начальные переменные процесса в виде
JSON-строки (также принимаютсяMapили@Json-объект) - Отправить команду и заблокироваться до подтверждения от
Zeebe(используйте возвращаемыйCompletionStageдля неблокирующего вызова)
Тот же клиент публикует сообщения (client.newPublishMessageCommand()) и выполняет любые другие команды Zeebe.
Настройка клиента¶
ZeebeClient можно донастроить необязательными компонентами графа, которые модуль подхватывает автоматически:
CredentialsProvider— авторизация дляZeebe(Camunda 8 SaaSили self-managed сOAuth);JsonMapper— пользовательскийJSON-маппер, используемыйZeebeClientдля (де)сериализации переменных;ScheduledExecutorService— пул потоков, используемый исполнителями заданий;ClientInterceptor—gRPC-перехватчик, применяемый к каналуZeebe(собираются все зарегистрированные перехватчики).
Например, чтобы аутентифицироваться в Camunda 8 через OAuth, предоставьте компонент CredentialsProvider:
@Module
public interface ZeebeAuthModule {
default CredentialsProvider zeebeCredentialsProvider() {
return CredentialsProvider.newCredentialsProviderBuilder()
.clientId("client-id")
.clientSecret("client-secret")
.audience("zeebe.camunda.io")
.authorizationServerUrl("https://login.cloud.camunda.io/oauth/token")
.build();
}
}
@Module
interface ZeebeAuthModule {
fun zeebeCredentialsProvider(): CredentialsProvider =
CredentialsProvider.newCredentialsProviderBuilder()
.clientId("client-id")
.clientSecret("client-secret")
.audience("zeebe.camunda.io")
.authorizationServerUrl("https://login.cloud.camunda.io/oauth/token")
.build()
}
Исполнители¶
Исполнитель — это обработчик, способный выполнять определенное задание в процессе.
Когда в процессе появляется задание нужного типа, Zeebe активирует его и передает одному из исполнителей.
Конфигурация¶
Существует конфигурация по умолчанию, которая применяется ко всем исполнителям при создании, а затем поверх нее
применяются именованные настройки конкретного исполнителя по типу исполнителя (Type).
Чтобы изменить настройки сразу для всех исполнителей, переопределите секцию default.
Чтобы изменить настройки только одного исполнителя, добавьте секцию с именем его типа, указанного в @JobWorker.
Если секция zeebe.worker.job не указана, используется встроенная конфигурация по умолчанию.
Пример полной конфигурации исполнителя, описанной в классе ZeebeWorkerConfig (указаны примеры значений или значения по умолчанию):
zeebe {
worker {
job {
default { //(1)!
enabled = true //(2)!
timeout = "15m" //(3)!
maxJobsActive = 32 //(4)!
requestTimeout = "15s" //(5)!
pollInterval = "100ms" //(6)!
tenantIds = [] //(7)!
streamEnabled = false //(8)!
streamTimeout = "15s" //(9)!
backoff {
minDelay = "100ms" //(11)!
maxDelay = "500ms" //(12)!
factor = 1.0 //(10)!
jitter = 1.1 //(13)!
}
}
}
}
}
- Тип исполнителя (
Type) или имя настроек по умолчаниюdefault - Включен ли исполнитель (по умолчанию:
true) - Максимальное время выполнения одного задания исполнителем (по умолчанию:
15m) - Максимальное количество заданий, которые будут одновременно активированы для этого исполнителя; используется для согласования скорости получения заданий со скоростью их обработки (
backpressure) (по умолчанию:32) - Ограничение времени запроса, который используется для опроса нового задания исполнителем (по умолчанию:
15s) - Максимальный интервал между опросами новых заданий; если после завершения работы новые задания не активированы, исполнитель периодически опрашивает брокер (по умолчанию:
100ms) - Идентификаторы
tenant, для которых исполнитель может получать задания (по умолчанию:[]) - Использовать ли потоковую передачу вместе с опросом для активации заданий (по умолчанию:
false) - Максимальное время жизни потока, если потоковая передача включена (по умолчанию:
15s) - Минимальная задержка повтора; из-за
jitterфактическая задержка может оказаться ниже этого минимума (по умолчанию:100ms) - Максимальная задержка повтора; из-за
jitterфактическая задержка может превысить это значение (по умолчанию:500ms) - Коэффициент умножения задержки: предыдущая задержка умножается на это значение (по умолчанию:
1.0) - Коэффициент
jitter: следующая задержка случайно изменяется в диапазоне+/-этого коэффициента (по умолчанию:1.1)
zeebe:
worker:
job:
default: #(1)!
enabled: true #(2)!
timeout: "15m" #(3)!
maxJobsActive: 32 #(4)!
requestTimeout: "15s" #(5)!
pollInterval: "100ms" #(6)!
tenantIds: [] #(7)!
streamEnabled: false #(8)!
streamTimeout: "15s" #(9)!
backoff:
minDelay: "100ms" #(11)!
maxDelay: "500ms" #(12)!
factor: 1.0 #(10)!
jitter: 1.1 #(13)!
- Тип исполнителя (
Type) или имя настроек по умолчаниюdefault - Включен ли исполнитель (по умолчанию:
true) - Максимальное время выполнения одного задания исполнителем (по умолчанию:
15m) - Максимальное количество заданий, которые будут одновременно активированы для этого исполнителя; используется для согласования скорости получения заданий со скоростью их обработки (
backpressure) (по умолчанию:32) - Ограничение времени запроса, который используется для опроса нового задания исполнителем (по умолчанию:
15s) - Максимальный интервал между опросами новых заданий; если после завершения работы новые задания не активированы, исполнитель периодически опрашивает брокер (по умолчанию:
100ms) - Идентификаторы
tenant, для которых исполнитель может получать задания (по умолчанию:[]) - Использовать ли потоковую передачу вместе с опросом для активации заданий (по умолчанию:
false) - Максимальное время жизни потока, если потоковая передача включена (по умолчанию:
15s) - Минимальная задержка повтора; из-за
jitterфактическая задержка может оказаться ниже этого минимума (по умолчанию:100ms) - Максимальная задержка повтора; из-за
jitterфактическая задержка может превысить это значение (по умолчанию:500ms) - Коэффициент умножения задержки: предыдущая задержка умножается на это значение (по умолчанию:
1.0) - Коэффициент
jitter: следующая задержка случайно изменяется в диапазоне+/-этого коэффициента (по умолчанию:1.1)
Чтобы переопределить настройки для одного исполнителя, добавьте секцию с ключом по типу исполнителя (Type),
указанному в @JobWorker. Именованная секция накладывается поверх default, который, в свою очередь, накладывается
поверх встроенных значений по умолчанию, поэтому в именованной секции достаточно перечислить только изменяемые ключи.
Установка enabled = false для именованного типа отключает только этого одного исполнителя.
zeebe {
worker {
job {
foo { //(1)!
timeout = "30s"
maxJobsActive = 8
}
bar { //(2)!
enabled = false
}
}
}
}
- Переопределяет только
timeoutиmaxJobsActiveдля исполнителя@JobWorker("foo"); все остальные настройки берутся изdefault - Отключает исполнителя
@JobWorker("bar"), не затрагивая остальную конфигурацию
Декларативные¶
Можно декларативно создавать исполнителей, которые будут
выполнять работу в рамках оркестратора Zeebe.
В аннотации @JobWorker указывается тип исполнителя (Type)
из процесса. По этому значению Zeebe связывает задание из BPMN-процесса с обработчиком в приложении.
Метод исполнителя может объявлять только параметры @JobVariable, @JobVariables и JobContext — любой другой тип
параметра отклоняется на этапе компиляции. Сырые JobClient и ActivatedJob доступны только в императивном исполнителе.
Параметр контекст¶
Можно внедрить контекст задания как аргумент метода.
JobContext содержит метаданные текущего задания, исполнителя и процесса.
JobContext предоставляет следующие методы только для чтения:
| Метод | Описание |
|---|---|
jobKey() |
Уникальный ключ активированного задания |
jobName() |
Имя/тип исполнителя, под которым зарегистрирован этот обработчик (значение @JobWorker) |
jobType() |
Тип активированного задания, как определено в BPMN-процессе |
jobWorker() |
Имя исполнителя, активировавшего задание на стороне брокера |
tenantId() |
Идентификатор tenant, которому принадлежит задание |
processId() |
Идентификатор BPMN-процесса |
processInstanceKey() |
Ключ экземпляра процесса, которому принадлежит задание |
processDefinitionVersion() |
Версия развернутого определения процесса |
processDefinitionKey() |
Ключ развернутого определения процесса |
elementId() |
Идентификатор BPMN-элемента, для которого создано задание |
elementInstanceKey() |
Ключ экземпляра BPMN-элемента |
headers() |
Пользовательские заголовки, заданные для задания в BPMN-модели |
retryCount() |
Количество оставшихся повторов для задания |
deadline() |
Момент (Instant), до которого задание эксклюзивно закреплено за исполнителем |
deadlineAsMillis() |
Тот же крайний срок, выраженный в миллисекундах эпохи |
variablesAsString() |
Сырые переменные задания в виде JSON-строки |
Параметр переменная¶
Можно внедрять переменные процесса как аргументы метода. Переменная процесса является частью состояния процесса и может быть установлена при старте процесса или как часть результата исполнителя.
Если указана хотя бы одна переменная через @JobVariable, сгенерированный исполнитель будет запрашивать у Zeebe
только такие переменные. Если @JobVariable не используется, исполнитель запрашивает все переменные задания.
Можно указать имя переменной явно в @JobVariable, либо будет использовано имя аргумента метода по умолчанию.
Так как переменные процесса передаются как JSON, аргумент метода может быть пользовательским типом, для которого
доступны JsonReader и JsonWriter.
Параметр переменные¶
Можно внедрить сразу несколько переменных процесса одним
аргументом метода через @JobVariables. Такой аргумент представляет все переменные задания как один JSON-объект.
Результат¶
Можно не только выполнять работу, но и возвращать результат как переменные в контекст процесса.
Результат можно возвращать как Map<String, Object>, который описывает структуру JSON-ответа.
Также можно возвращать именованный результат как одну переменную. Это аналог одного ключа и значения в объекте
Map<String, Object>.
В таком случае обязательно требуется указать имя переменной в аннотации @JobVariable:
Ошибки¶
Если требуется завершить исполнение ошибкой процесса, бросьте JobWorkerException.
В исключении можно указать код ошибки, сообщение и переменные процесса, если они нужны.
Это исключение преобразуется в команду throwError для Zeebe: getCode(), сообщение и getVariables() исключения
передаются как код ошибки, сообщение об ошибке и переменные команды.
Если обработчик выбрасывает любое другое исключение, модуль оборачивает его в JobWorkerException с одним из следующих
встроенных кодов:
| Код | Когда используется |
|---|---|
DESERIALIZATION |
Переменную задания не удалось прочитать/десериализовать в аргумент метода |
SERIALIZATION |
Результат исполнителя не удалось записать/сериализовать в переменные |
UNEXPECTED |
Из синхронного обработчика было выброшено неожиданное исключение |
INTERNAL |
Резервный код для любой другой ошибки, не охваченной выше |
Императивные¶
Можно также создавать более низкоуровневые исполнители и напрямую работать с контрактами ZeebeClient.
Для этого компонент должен реализовать интерфейс KoraJobWorker.
@Component
public final class SomeJob implements KoraJobWorker {
@Override
public String type() {
return "someJobType";
}
@Override
public List<String> fetchVariables() {
return List.of("startId"); //(1)!
}
@Override
public CompletionStage<FinalCommandStep<?>> handle(JobClient client, ActivatedJob job) {
return CompletableFuture.completedFuture(client.newCompleteCommand(job));
}
}
- Из
Zeebeзапрашиваются только эти переменные; верните пустой список (по умолчанию), чтобы запросить все переменные
@Component
class SomeJob : KoraJobWorker {
override fun type(): String = "someJobType"
override fun fetchVariables(): List<String> = listOf("startId") //(1)!
override fun handle(client: JobClient, job: ActivatedJob): CompletionStage<FinalCommandStep<*>> {
return CompletableFuture.completedFuture(client.newCompleteCommand(job))
}
}
- Из
Zeebeзапрашиваются только эти переменные; верните пустой список (по умолчанию), чтобы запросить все переменные
Метод fetchVariables() — императивный аналог @JobVariable: он определяет, какие переменные процесса Zeebe
отправляет вместе с заданием. По умолчанию он возвращает пустой список, что запрашивает все переменные; непустой список
ограничивает передаваемые данные только этими переменными. В отличие от декларативных исполнителей, handle получает
сырые JobClient и ActivatedJob и отвечает за завершение задания (например, через client.newCompleteCommand(job)).
Сигнатуры¶
Доступные сигнатуры для методов исполнителя из коробки:
Под T подразумевается тип возвращаемого значения, либо Void.
Если результат равен null или Optional.empty(), задание будет завершено без добавления переменных.
T myMethod()Optional<T> myMethod()CompletionStage<T> myMethod()CompletionStageMono<T> myMethod()Project Reactor (надо подключить зависимость)
Под T подразумевается тип возвращаемого значения, либо T?, либо Unit.
Если результат равен null, задание будет завершено без добавления переменных.
myMethod(): TmyMethod(): Deferred<T>Kotlin Coroutine (надо подключить зависимость какimplementation)