Kora облачно ориентированный серверный фреймворк написанный на Java для написания Java / Kotlin приложений с упором на производительность, эффективность, прозрачность сделанный выходцами из Т-Банк / Тинькофф

Kora is a cloud-oriented server-side Java framework for writing Java / Kotlin applications with a focus on performance, efficiency and transparency

Перейти к содержанию

Camunda Zeebe

Экспериментальный модуль

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

Модуль подключает клиент Camunda 8 (Zeebe) и создает исполнителей заданий для внешнего оркестратора процессов. В Kora такой исполнитель объявляется обычным компонентом: метод с аннотацией @JobWorker получает переменные процесса, выполняет работу и возвращает результат, который будет передан обратно в Zeebe.

Подключение

Зависимость build.gradle:

implementation "ru.tinkoff.kora.experimental:camunda-zeebe-worker"

Модуль:

@KoraApp
public interface Application extends ZeebeWorkerModule { }

Зависимость build.gradle.kts:

implementation("ru.tinkoff.kora.experimental:camunda-zeebe-worker")

Модуль:

@KoraApp
interface Application : ZeebeWorkerModule

Конфигурация

Пример полной конфигурации клиента, описанной в классе 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"
                }
            }
        }
    }
}
  1. Максимальное количество потоков для исполнителей заданий (по умолчанию: количество ядер процессора, но не меньше 2)
  2. Время без активности чтения перед отправкой проверки KeepAlive (по умолчанию: 45s)
  3. Использовать ли TLS при подключении (по умолчанию: true)
  4. Файловый путь до сертификата для подключения; если не указан, используется системный сертификат (по умолчанию не указано, необязательно)
  5. Максимальное время ожидания проверки доступности топологии при запуске клиента (по умолчанию не указано, необязательно)
  6. URL для подключения по gRPC (обязательная, по умолчанию не указано)
  7. Время, в течение которого сообщение должно храниться на брокере при отправке по gRPC (по умолчанию: 1h)
  8. Максимальный размер входящего сообщения по gRPC (по умолчанию: 4Mib)
  9. Включена ли политика повторов для gRPC-соединения (по умолчанию: true)
  10. Количество попыток (по умолчанию: 5)
  11. Начальная задержка между попытками (по умолчанию: 100ms)
  12. Максимальная задержка между попытками (по умолчанию: 5s)
  13. Коэффициент увеличения задержки между попытками (по умолчанию: 3.0)
  14. URL для подключения к REST-адресу Zeebe; если указан, клиент предпочитает REST вместо gRPC для поддерживаемых операций (обязательная внутри необязательной секции rest, по умолчанию не указано)
  15. Пути для поиска ресурсов, которые будут загружены в оркестратор после запуска (по умолчанию: [])
  16. Максимальное время ожидания загрузки ресурсов (по умолчанию: 45s)
  17. Включает логирование модуля (по умолчанию: false)
  18. Включает метрики модуля (по умолчанию: true)
  19. Настройка SLO для метрик (по умолчанию: ru.tinkoff.kora.telemetry.common.TelemetryConfig.MetricsConfig#DEFAULT_SLO)
  20. Настройка тегов для метрик (по умолчанию: {})
  21. Включает трассировку модуля (по умолчанию: true)
  22. Настройка атрибутов для трассировки (по умолчанию: {})
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
  1. Максимальное количество потоков для исполнителей заданий (по умолчанию: количество ядер процессора, но не меньше 2)
  2. Время без активности чтения перед отправкой проверки KeepAlive (по умолчанию: 45s)
  3. Использовать ли TLS при подключении (по умолчанию: true)
  4. Файловый путь до сертификата для подключения; если не указан, используется системный сертификат (по умолчанию не указано, необязательно)
  5. Максимальное время ожидания проверки доступности топологии при запуске клиента (по умолчанию не указано, необязательно)
  6. URL для подключения по gRPC (обязательная, по умолчанию не указано)
  7. Время, в течение которого сообщение должно храниться на брокере при отправке по gRPC (по умолчанию: 1h)
  8. Максимальный размер входящего сообщения по gRPC (по умолчанию: 4Mib)
  9. Включена ли политика повторов для gRPC-соединения (по умолчанию: true)
  10. Количество попыток (по умолчанию: 5)
  11. Начальная задержка между попытками (по умолчанию: 100ms)
  12. Максимальная задержка между попытками (по умолчанию: 5s)
  13. Коэффициент увеличения задержки между попытками (по умолчанию: 3.0)
  14. URL для подключения к REST-адресу Zeebe; если указан, клиент предпочитает REST вместо gRPC для поддерживаемых операций (обязательная внутри необязательной секции rest, по умолчанию не указано)
  15. Пути для поиска ресурсов, которые будут загружены в оркестратор после запуска (по умолчанию: [])
  16. Максимальное время ожидания загрузки ресурсов (по умолчанию: 45s)
  17. Включает логирование модуля (по умолчанию: false)
  18. Включает метрики модуля (по умолчанию: true)
  19. Настройка SLO для метрик (по умолчанию: ru.tinkoff.kora.telemetry.common.TelemetryConfig.MetricsConfig#DEFAULT_SLO)
  20. Настройка тегов для метрик (по умолчанию: {})
  21. Включает трассировку модуля (по умолчанию: true)
  22. Настройка атрибутов для трассировки (по умолчанию: {})

Предоставляемые метрики модуля описаны в разделе Справочник метрик.

Развертывание ресурсов

Если в deployment.resources указаны пути, модуль во время запуска находит ресурсы в classpath и развертывает их в Zeebe через компонент ZeebeResourceDeployment. Развертываются как BPMN-процессы, так и DMN-решения, найденные в настроенных расположениях. Поддерживаются только пути с префиксом classpath:, например classpath:bpm; другие расположения логируются и пропускаются.

Разместите развертываемые ресурсы в соответствующей директории classpath:

src/main/resources/
└── bpm/
    └── demo.bpmn
zeebe {
    client {
        deployment {
            resources = "classpath:bpm" //(1)!
        }
    }
}
  1. Одно или несколько расположений в classpath для поиска BPMN / DMN-ресурсов (одно значение или список)
zeebe:
  client:
    deployment:
      resources: "classpath:bpm" #(1)!
  1. Одно или несколько расположений в 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)!
    }
}
  1. Идентификатор BPMN-процесса, который нужно запустить
  2. Запуск последней развернутой версии процесса
  3. Начальные переменные процесса в виде JSON-строки (также принимаются Map или @Json-объект)
  4. Отправить команду и заблокироваться до подтверждения от 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)!
    }
}
  1. Идентификатор BPMN-процесса, который нужно запустить
  2. Запуск последней развернутой версии процесса
  3. Начальные переменные процесса в виде JSON-строки (также принимаются Map или @Json-объект)
  4. Отправить команду и заблокироваться до подтверждения от Zeebe (используйте возвращаемый CompletionStage для неблокирующего вызова)

Тот же клиент публикует сообщения (client.newPublishMessageCommand()) и выполняет любые другие команды Zeebe.

Настройка клиента

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

  • CredentialsProvider — авторизация для Zeebe (Camunda 8 SaaS или self-managed с OAuth);
  • JsonMapper — пользовательский JSON-маппер, используемый ZeebeClient для (де)сериализации переменных;
  • ScheduledExecutorService — пул потоков, используемый исполнителями заданий;
  • ClientInterceptorgRPC-перехватчик, применяемый к каналу 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)!
                }
            }
        }
    }
}
  1. Тип исполнителя (Type) или имя настроек по умолчанию default
  2. Включен ли исполнитель (по умолчанию: true)
  3. Максимальное время выполнения одного задания исполнителем (по умолчанию: 15m)
  4. Максимальное количество заданий, которые будут одновременно активированы для этого исполнителя; используется для согласования скорости получения заданий со скоростью их обработки (backpressure) (по умолчанию: 32)
  5. Ограничение времени запроса, который используется для опроса нового задания исполнителем (по умолчанию: 15s)
  6. Максимальный интервал между опросами новых заданий; если после завершения работы новые задания не активированы, исполнитель периодически опрашивает брокер (по умолчанию: 100ms)
  7. Идентификаторы tenant, для которых исполнитель может получать задания (по умолчанию: [])
  8. Использовать ли потоковую передачу вместе с опросом для активации заданий (по умолчанию: false)
  9. Максимальное время жизни потока, если потоковая передача включена (по умолчанию: 15s)
  10. Минимальная задержка повтора; из-за jitter фактическая задержка может оказаться ниже этого минимума (по умолчанию: 100ms)
  11. Максимальная задержка повтора; из-за jitter фактическая задержка может превысить это значение (по умолчанию: 500ms)
  12. Коэффициент умножения задержки: предыдущая задержка умножается на это значение (по умолчанию: 1.0)
  13. Коэффициент 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)!
  1. Тип исполнителя (Type) или имя настроек по умолчанию default
  2. Включен ли исполнитель (по умолчанию: true)
  3. Максимальное время выполнения одного задания исполнителем (по умолчанию: 15m)
  4. Максимальное количество заданий, которые будут одновременно активированы для этого исполнителя; используется для согласования скорости получения заданий со скоростью их обработки (backpressure) (по умолчанию: 32)
  5. Ограничение времени запроса, который используется для опроса нового задания исполнителем (по умолчанию: 15s)
  6. Максимальный интервал между опросами новых заданий; если после завершения работы новые задания не активированы, исполнитель периодически опрашивает брокер (по умолчанию: 100ms)
  7. Идентификаторы tenant, для которых исполнитель может получать задания (по умолчанию: [])
  8. Использовать ли потоковую передачу вместе с опросом для активации заданий (по умолчанию: false)
  9. Максимальное время жизни потока, если потоковая передача включена (по умолчанию: 15s)
  10. Минимальная задержка повтора; из-за jitter фактическая задержка может оказаться ниже этого минимума (по умолчанию: 100ms)
  11. Максимальная задержка повтора; из-за jitter фактическая задержка может превысить это значение (по умолчанию: 500ms)
  12. Коэффициент умножения задержки: предыдущая задержка умножается на это значение (по умолчанию: 1.0)
  13. Коэффициент jitter: следующая задержка случайно изменяется в диапазоне +/- этого коэффициента (по умолчанию: 1.1)

Чтобы переопределить настройки для одного исполнителя, добавьте секцию с ключом по типу исполнителя (Type), указанному в @JobWorker. Именованная секция накладывается поверх default, который, в свою очередь, накладывается поверх встроенных значений по умолчанию, поэтому в именованной секции достаточно перечислить только изменяемые ключи. Установка enabled = false для именованного типа отключает только этого одного исполнителя.

zeebe {
    worker {
        job {
            foo { //(1)!
                timeout = "30s"
                maxJobsActive = 8
            }
            bar { //(2)!
                enabled = false
            }
        }
    }
}
  1. Переопределяет только timeout и maxJobsActive для исполнителя @JobWorker("foo"); все остальные настройки берутся из default
  2. Отключает исполнителя @JobWorker("bar"), не затрагивая остальную конфигурацию
zeebe:
  worker:
    job:
      foo: #(1)!
        timeout: "30s"
        maxJobsActive: 8
      bar: #(2)!
        enabled: false
  1. Переопределяет только timeout и maxJobsActive для исполнителя @JobWorker("foo"); все остальные настройки берутся из default
  2. Отключает исполнителя @JobWorker("bar"), не затрагивая остальную конфигурацию

Декларативные

Можно декларативно создавать исполнителей, которые будут выполнять работу в рамках оркестратора Zeebe.

В аннотации @JobWorker указывается тип исполнителя (Type) из процесса. По этому значению Zeebe связывает задание из BPMN-процесса с обработчиком в приложении.

Метод исполнителя может объявлять только параметры @JobVariable, @JobVariables и JobContext — любой другой тип параметра отклоняется на этапе компиляции. Сырые JobClient и ActivatedJob доступны только в императивном исполнителе.

@Component
public final class SomeJob {

    @JobWorker("someJobType")
    public void process() {
        // do something
    }
}
@Component
class SomeJob {

    @JobWorker("someJobType")
    fun process() {
        // do something
    }
}

Параметр контекст

Можно внедрить контекст задания как аргумент метода. JobContext содержит метаданные текущего задания, исполнителя и процесса.

@Component
public final class SomeJob {

    @JobWorker("someJobType")
    public void process(JobContext context) {
        // do something
    }
}
@Component
class SomeJob {

    @JobWorker("someJobType")
    fun process(context: JobContext) {
        // do something
    }
}

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-строки
@Component
public final class SomeJob {

    @JobWorker("someJobType")
    public void process(JobContext context) {
        logger.info("Job {} of process {} at element {} with deadline {}",
                context.jobType(), context.processInstanceKey(), context.elementId(), context.deadline());
    }
}
@Component
class SomeJob {

    @JobWorker("someJobType")
    fun process(context: JobContext) {
        logger.info("Job {} of process {} at element {} with deadline {}",
            context.jobType(), context.processInstanceKey(), context.elementId(), context.deadline())
    }
}

Параметр переменная

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

Если указана хотя бы одна переменная через @JobVariable, сгенерированный исполнитель будет запрашивать у Zeebe только такие переменные. Если @JobVariable не используется, исполнитель запрашивает все переменные задания.

@Component
public final class SomeJob {

    @JobWorker("someJobType")
    public void process(@JobVariable("startId") String id) {
        // do something
    }
}
@Component
class SomeJob {

    @JobWorker("someJobType")
    fun process(@JobVariable("startId") id: String) {
        // do something
    }
}

Можно указать имя переменной явно в @JobVariable, либо будет использовано имя аргумента метода по умолчанию.

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

@Component
public final class SomeJob {

    @Json
    public record User(String name, int code) { }

    @JobWorker("someJobType")
    public void process(@JobVariable User user) {
        // do something
    }
}
@Component
class SomeJob {

    data class User(val name: String, val code: Int)

    @JobWorker("someJobType")
    fun process(@JobVariable user: User) {
        // do something
    }
}

Параметр переменные

Можно внедрить сразу несколько переменных процесса одним аргументом метода через @JobVariables. Такой аргумент представляет все переменные задания как один JSON-объект.

@Component
public final class SomeJob {

    @Json
    public record User(String name, int code) { }

    @Json
    public record UserContext(String startId, User user) { }

    @JobWorker("someJobType")
    public void process(@JobVariables UserContext userContext) {
        // do something
    }
}
@Component
class SomeJob {

    data class User(val name: String, val code: Int)

    data class UserContext(val startId: String, val user: User)

    @JobWorker("someJobType")
    fun process(@JobVariables userContext: UserContext) {
        // do something
    }
}

Результат

Можно не только выполнять работу, но и возвращать результат как переменные в контекст процесса.

Результат можно возвращать как Map<String, Object>, который описывает структуру JSON-ответа.

@Component
public final class SomeJob {

    @JobWorker("someJobType")
    public Map<String, Object> process() {
        // do something
    }
}
@Component
class SomeJob {

    @JobWorker("someJobType")
    fun process(): Map<String, Any> {
        // do something
    }
}

Также можно возвращать именованный результат как одну переменную. Это аналог одного ключа и значения в объекте Map<String, Object>.

В таком случае обязательно требуется указать имя переменной в аннотации @JobVariable:

@Component
public final class SomeJob {

    @Json
    public record User(String name, int code) { }

    @JobVariable("user")
    @JobWorker("someJobType")
    public User process() {
        // do something
    }
}
@Component
class SomeJob {

    data class User(val name: String, val code: Int)

    @JobVariable("user")
    @JobWorker("someJobType")
    fun process(): User {
        // do something
    }
}

Ошибки

Если требуется завершить исполнение ошибкой процесса, бросьте JobWorkerException. В исключении можно указать код ошибки, сообщение и переменные процесса, если они нужны. Это исключение преобразуется в команду throwError для Zeebe: getCode(), сообщение и getVariables() исключения передаются как код ошибки, сообщение об ошибке и переменные команды.

Если обработчик выбрасывает любое другое исключение, модуль оборачивает его в JobWorkerException с одним из следующих встроенных кодов:

Код Когда используется
DESERIALIZATION Переменную задания не удалось прочитать/десериализовать в аргумент метода
SERIALIZATION Результат исполнителя не удалось записать/сериализовать в переменные
UNEXPECTED Из синхронного обработчика было выброшено неожиданное исключение
INTERNAL Резервный код для любой другой ошибки, не охваченной выше
@Component
public final class SomeJob {

    @JobWorker("someJobType")
    public User process() {
        throw new JobWorkerException("DOESNT_WORK"); //(1)!
    }
}
  1. Дополнительные перегрузки принимают сообщение/причину и Map<String, Object> переменных для добавления в команду throwError
@Component
class SomeJob {

    @JobWorker("someJobType")
    fun process(): User {
        throw JobWorkerException("DOESNT_WORK") //(1)!
    }
}
  1. Дополнительные перегрузки принимают сообщение/причину и Map<String, Any> переменных для добавления в команду throwError

Императивные

Можно также создавать более низкоуровневые исполнители и напрямую работать с контрактами 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));
    }
}
  1. Из 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))
    }
}
  1. Из Zeebe запрашиваются только эти переменные; верните пустой список (по умолчанию), чтобы запросить все переменные

Метод fetchVariables() — императивный аналог @JobVariable: он определяет, какие переменные процесса Zeebe отправляет вместе с заданием. По умолчанию он возвращает пустой список, что запрашивает все переменные; непустой список ограничивает передаваемые данные только этими переменными. В отличие от декларативных исполнителей, handle получает сырые JobClient и ActivatedJob и отвечает за завершение задания (например, через client.newCompleteCommand(job)).

Сигнатуры

Доступные сигнатуры для методов исполнителя из коробки:

Под T подразумевается тип возвращаемого значения, либо Void. Если результат равен null или Optional.empty(), задание будет завершено без добавления переменных.

Под T подразумевается тип возвращаемого значения, либо T?, либо Unit. Если результат равен null, задание будет завершено без добавления переменных.