Выполнение правил трансформации события#
Описание#
Правила трансформации событий описываются в файле *.tr. По умолчанию обработчик не производит никаких действий над потоком, все необходимые действия должны указываться явно.
При выполнении правил трансформации события используются идентификаторы для входа и выхода. Описываются зарезервированными словами:
INPUT, Input, input — для адресации элемента из входящего сообщения;
OUTPUT, Output, output — для адресации элемента из выходящего сообщения.
Или
IN, In, in — для адресации элемента из входящего сообщения;
OUT, Out, out — для адресации элемента из выходящего сообщения.
Вход и выход имеют одинаковую структуру:
headers — заголовки сообщения;
key — ключ сообщения Kafka;
body — тело сообщения;
sourceTopic — топик источника;
targetTopic — топик назначения;
branch — имя ветки для маршрутизации.
Так же в этих корнях могут быть массивы элементов в случае нескольких сообщений на вход или выход.
При использовании INPUT и OUTPUT после (возможно, после индекса) обязательно должен быть указан тип узла:
headers;
key;
body;
sourceTopic;
targetTopic;
branch.
Примеры использования#
Пример правил трансформации:
out = in
out.branch = "second"
headers[>].TargetTopic = "output"
out[>].first = in
headers[+>].BranchName = "transformSecond"
out[+>].second = in
headers[+>].BranchName = "transformThird"
out[+>].third = in
Пример правил трансформации с указанием типа узла:
output.body = in
output.body.branch = "second"
output[>].targetTopic = "output"
output[>].body.first = in
output[+>].branch = "transformSecond"
output[>].body.second = in
output[+>].branch = "transformThird"
output[>].body.third = in
OUTPUT.body.name.var = INPUT.body.text
Использование правил трансформации при обратной совместимости#
Для обратной совместимости применимы старые методы обращения к входам и выходам.
При трансформации события происходит выполнение DSL правил трансформации:
Вход
INсодержит входящее событие;Environment.SourceTopicсодержит имя топика, откуда было прочитано сообщение;Environment.SourceKeyсодержит значение ключа из сообщения Apache Kafka в виде строки;Environment.kafka.headersсодержит словарь заголовков из сообщения Apache Kafka;Environment.cacheсодержит элементы, сохраняющиеся между вызовами правил трансформации;Environment.loadedVariablesсодержит элементы, загруженные из JSON файла из параметраenvironmentVariablesшага трансформации событий.
Выход OUT должен содержать исходящее событие или массив исходящих событий.
Выход HEADERS может содержать следующие поля:
TargetTopicимя топика назначения события;TargetKeyзначение ключа в сообщении Apache Kafka;kafkaсловарь заголовков сообщения Apache Kafka, не поддерживает вложенные словари;BranchNameимя следующего шага из побочной ветки потока.
Если исходящих событий больше одного, то выход HEADERS должен содержать массив элементов с перечисленными выше полями. Порядок элементов должен соответствовать порядку событий в выходе OUT.
При работе с INPUT и OUTPUT вместо выхода HEADERS нужно использовать соответствующие значения узлов INPUT и OUTPUT.
Поддерживается вызов функции:
return("decline"), который прерывает функцию трансформации и фильтрует полученное событие;return("passthrough"), который передает входящие сообщения на следующий шаг без изменений (имеется возможность изменять заголовки и другую информацию при помощи DSL).
Вызов функции return() с любым другим параметром прервет дальнейшее выполнение правил трансформации, при этом уже заполненные выходы OUT и HEADERS будут обработаны и преобразованы в исходящие события.
Функция агрегации событий#
Правила функции агрегации событий аналогичны правилам трансформации события, только функция агрегации принимает на вход массив событий, выполняет правила трансформации и на выход возвращает одно или несколько событий в зависимости от логики правил трансформации.
При выполнении DSL правил трансформации:
Вход
INсодержит события из окна агрегации;Environment.sourcesсодержит массив переменных окружения для каждого события, индекс элементов соответствует индексам событий изIN. Каждый элемент массива содержит поля:SourceTopic— содержит имя топика, откуда было прочитано сообщение;SourceKey— содержит значение ключа из сообщения Apache Kafka в виде строки;kafka.headers— содержит словарь заголовков из сообщения Apache Kafka;
Environment.cache— содержит элементы, сохраняющиеся между вызовами функции агрегации;Environment.loadedVariables— содержит элементы, загруженные из JSON файла из параметраenvironmentVariablesшага агрегации;Environment.AggregationKey— содержит значение ключа агрегации;Environment.AggregationWindowTimeMs— содержит продолжительность существования окна агрегации в миллисекундах;Environment.AggregationEventCount— содержит количество событий, переданных в функцию агрегации.