Трансформация событий#
Шаг предназначен для преобразования события из одного формата в другой с помощью модуля declarative-mapper.
Представляет собой объект с полем type со значением "dsl" и полями:
name— имя шага;dsl— описание расположения правил трансформации;environmentVariables— загрузка переменных из файла;format— описание форматов входящего и исходящего событий:input— формат входящего события;output— формат исходящего события;
stateFull— сохранениеEnvironment.cacheв состоянииFlink, логическое, по умолчаниюfalse;destination— следующий шаг потока;error— шаг обработки в случае возникновения ошибки;database— настройки для базы, используемой в правилах трансформации;async— признак асинхронного выполнения, по умолчаниюfalse;threadPoolSize— размер пула потоков для асинхронного выполнения трансформации;timeout— тайм-аут выполнения трансформации;key— вычисление ключа для определения дубля события;deduplicator.duration— продолжительность времени для отсечения последующих сообщений;shouldLogDebug— флаг для вывода информации о результатах выполнения шагов dsl.
Поле deduplicator.duration должно иметь значение вида — <величина> <единица измерения>. Например: 2 milliseconds.
Единица измерения может принимать следующие значения:
по умолчанию — миллисекунда;
ns, nano, nanos, nanosecond, nanoseconds;
us, micro, micros, microsecond, microseconds;
ms, milli, millis, millisecond, milliseconds;
s, second, seconds;
m, minute, minutes;
h, hour, hours;
d, day, days.
Значения полей dsl и environmentVariables представляет собой объект с полями:
source— тип источника правил, принимает значения:file— правила расположены в файловой системе;classpath— правила расположены в classpath JVM;
path— путь до файла на диске или в classpath JVM.
Значения полей input и output могут быть следующие:
json— для событий в формате JSON;xml— для событий в формате XML;событие в формате
csv:type— со значениемcsv;columnSeparator— разделитель полей события;fieldNames— список имен полей события для использования при трансформации, также определяет порядок полей.
объект с полями:
typeсо значениемavro;schema— список схем, значение аналогично полюdsl, объект имеющий логическое полеdefaultсо значениемtrue, будет считаться схемой по умолчанию.
Пример#
transformStep: {
name: "transformation-avro-json"
type: "dsl"
format: {
input: {
type: "avro"
schema: [{
source: "file"
path: "/file/path/to/avro/schema"
},{
source: "classpath"
path: "/path/to/default/avro/schema"
default: true
}]
}
output: "json"
}
dsl: {
source: "file"
path: "Through.tr"
}
environmentVariables: {
source: "file"
path: "/path/to/environment.json"
}
destination: {}
error: {}
}
AVRO-схема выбирается исходя из значения заголовка "schemaName" в сообщении Apache Kafka, если заголовок не задан,
то используется схема по умолчанию или первая в списке в случае, если никакая схема не помечена полем default со значением true.
Преобразование форматов событий в формате csv#
transformStep: {
name: "transformation-csv-csv"
type: "dsl"
format: {
input: {
type: "csv"
columnSeparator: ";"
fieldNames: ["field_1", "field_2"] # можно не указывать и при трансформации использовать порядковые номера полей
}
output: {
type: "csv"
columnSeparator: ","
fieldNames: ["field_1", "field_2", "field_3"] # в output указывать обязательно
}
}
dsl: {
source: "file"
path: "Through.tr"
}
destination: {}
}
Порядок имен полей в "filedNames" имеет значение — значения будут прочитаны/записаны в том порядке, в котором указаны в данном параметре.
Шаг фильтрации/маршрутизации#
Шаг предназначен для фильтрации/маршрутизации события.
Представляет собой объект с полем type со значением "filter" и полями:
name— имя шага;input— формат входящего события;mode— при значении «all» проверяться будут все условия, иначе до первого успешного;filter— массив маршрутизации;condition— условие на DSL, в результате дающее в итоге true или false;destination— следующий шаг при выполнении условия, если поле имеет значение «decline» - обработка прекращается, сообщение отбрасывается;output— формат исходящего события;log— лог при выполнении условия;level— уровень логирования в соответствии с функциями логирования dsl;message— сообщение на DSL (аналогично стандартным аргументам функции логирования DSL);eventId— eventId на DSL (аналогично стандартным аргументам функции логирования DSL);input— input на DSL (аналогично стандартным аргументам функции логирования DSL);output— output на DSL (аналогично стандартным аргументам функции логирования DSL);
default— маршрутизация по умолчанию;condition— условие на DSL, в результате дающее в итоге true или false;destination— следующий шаг при выполнении условия, если поле имеет значение «decline» - обработка прекращается, сообщение отбрасывается;output— формат исходящего события;log— лог при выполнении условия;level— уровень логирования в соответствии с функциями логирования dsl;message— сообщение на DSL (аналогично стандартным аргументам функции логирования DSL);eventId— eventId на DSL (аналогично стандартным аргументам функции логирования DSL);input— input на DSL (аналогично стандартным аргументам функции логирования DSL);output— output на DSL (аналогично стандартным аргументам функции логирования DSL).
Алгоритм работы шага фильтрации#
Алгоритм должен проверять условия блоков condition последовательно в зависимости от значения поля modе:
по умолчанию
modене указывается и работа происходит до первого выполнения условия. После первого сработавшего условия остальные условия не проверяем;если указан
modе = "all", то проверяем условия всех блоков и при выполнении условия - выполняем содержимое блока.
Если выполнено условие, для которого указан destination: «decline» - обработка данного сообщения прекращается, сообщение никуда не отправляется.
Если ни одно условие не выполнено - работаем по блоку default.
Блок default должен присутствовать обязательно. Если его нет - выдавать ошибку.
Примеры#
Пример с одним условием#
transformStep: {
name: "transformation"
type: "filter"
input: "json"
filter: [{
condition: "IN.message == \"test2\" or IN.message == \"test1\""
destination: ${transformThird}
output: "json"
}
],
default: {
destination: "decline"
output: "xml"
}
}
Пример с двумя условиями#
transformStep: {
name: "transformation"
type: "filter"
input: "json"
filter: [
{
condition: "IN.message == \"test1\""
destination: ${transformSecond}
output: "json"
log: {
level: "info"
message: "\"Сообщение с id \" || in.message || \" отправлено в топик Topic_Name\""
input: "\"testInput\""
output: "\"testOutput\""
}
},
{
condition: "IN.message == \"test2\" or IN.message == \"test1\""
destination: ${transformThird}
output: "json"
},
],
default: {
destination: ${destinationThird}
output: "xml"
}
}
Пример для формата «HOCON»#
{
flow: {
name: "JobName"
source: ${source}
}
source: {
name: "source"
type: "source"
topic: "input"
config: {
consumer: {
"group.id": "event-process-flow-group"
"client.id": "event-process-flow-client"
}
producer: {
"client.id": "event-process-flow-producer"
}
}
destination: ${transformStep}
}
transformStep: {
name: "transformation step"
type: "filter"
modе: "all"
filter: [
{
condition: "условие на DSL"
destination: ${destination1} # целевой шаг
format: "json" # формат выходного сообщения
log {
level: "info"
message: "Сообщение с id" || in.messageID || "отправлено в топик Topic_Name"
input: environment.input_topic # опциональное поле
output: "topic" # опциональное поле
}
},
{
condition: "условие на DSL"
destination: ${destination2} # целевой шаг
format: "xml" # формат выходного сообщения
},
{
condition: "условие на DSL"
destination: "decline" # при выполнении условия сообщение отбрасывается
}
]
default {
destination: ${destinationDefault}
format: "json"
}
}
destination1: {
name: "destination1"
type: "destination"
topic: "output1"
config: ${defaults.kafka}
}
destination2: {
name: "destination2"
type: "destination"
topic: "output2"
config: ${defaults.kafka}
}
destinationDefault: {
name: "destinationDefault"
type: "destination"
topic: "outputDefault"
config: ${defaults.kafka}
}
}
Использование режима all при проверке условий с разными форматами сообщений#
При использовании режима mode = all возможно указание нескольких валидных форматов сообщений с помощью блока additional для следующего случая:
сообщение может соответствовать нескольким условиям;
формат выходного сообщения в одном из блоков условия отличается от формата входного сообщения.
Пример конфигурации формата входного сообщения с дополнительными валидными форматами#
input: {
type: "json"
additional: [
{
type: "xml"
},
{
type: "avro"
}
]
}
Пример шага трансформации для формата «HOCON»#
transformStep: {
name: "transformation"
type: "filter"
modе: "all"
input: {
type: "json"
additional: [
{
type: "xml"
}
]
}
filter: [
{
condition: "IN.message == \"test1\""
destination: ${transformSecond}
output: "xml"
log: {
level: "info"
message: "\"Сообщение с id \" || in.message || \" отправлено в топик Topic_Name\""
input: "\"testInput\""
output: "\"testOutput\""
}
},
{
condition: "IN.message == \"test2\" or IN.message == \"test1\""
destination: ${transformThird}
output: "json"
}
]
default: {
destination: ${destinationThird}
output: "xml"
}
}