Трансформация событий#

Шаг предназначен для преобразования события из одного формата в другой с помощью модуля declarative-mapper.

Представляет собой объект с полем type со значением "dsl" и полями:

  1. name — имя шага;

  2. dsl — описание расположения правил трансформации;

  3. environmentVariables — загрузка переменных из файла;

  4. format — описание форматов входящего и исходящего событий:

    • input — формат входящего события;

    • output — формат исходящего события;

  5. stateFull — сохранение Environment.cache в состоянии Flink, логическое, по умолчанию false;

  6. destination — следующий шаг потока;

  7. error — шаг обработки в случае возникновения ошибки;

  8. database — настройки для базы, используемой в правилах трансформации;

  9. async — признак асинхронного выполнения, по умолчанию false;

  10. threadPoolSize — размер пула потоков для асинхронного выполнения трансформации;

  11. timeout — тайм-аут выполнения трансформации;

  12. key — вычисление ключа для определения дубля события;

  13. deduplicator.duration — продолжительность времени, в течение которого гарантируется отсутствие дублирующих сообщений;

  14. 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 представляет собой объект с полями:

  1. source — тип источника правил, принимает значения:

    • file — правила расположены в файловой системе;

    • classpath — правила расположены в classpath JVM;

  2. path — путь до файла на диске или в classpath JVM.

Значения полей input и output могут быть следующие:

  1. json — для событий в формате JSON;

  2. xml — для событий в формате XML;

  3. событие в формате csv:

    • type — со значением csv;

    • columnSeparator — разделитель полей события;

    • fieldNames — список имен полей события для использования при трансформации, также определяет порядок полей.

  4. объект с полями:

    • 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" имеет значение — значения будут прочитаны/записаны в том порядке, в котором указаны в данном параметре.

Формат CloudEvents для шагов преобразования форматов и агрегации событий (применимо только к компоненту EVPC)#

Пример преобразования форматов событий в формате CloudEvents

transformStep: {
  name: "transformation-cloud-event"
  type: "dsl"
  format: {
    input: "cloudEvent"
    output: "cloudEvent"
  }
  dsl: {
    source: "file"
    path: "Through.tr"
  }
  environmentVariables: {
    source: "file"
    path: "./environment.json"
  }
  destination: {}
  error: {}
}

В этом случае логика преобразования выбирается на основе атрибута datacontenttype. При отсутствии атрибута форматом данных события CloudEvent считается JSON (datacontenttype = "application/json"). Формат выходного события определяется также на основе атрибута datacontenttype.

Шаг фильтрации/маршрутизации#

Шаг предназначен для фильтрации/маршрутизации события.

Представляет собой объект с полем type со значением "filter" и полями:

  1. name — имя шага;

  2. input — формат входящего события;

  3. mode — при значении «all» проверяться будут все условия, иначе до первого успешного;

  4. 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);

  5. 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}
  }
}

Пример для формата «YAML»#

flow:
  name: "Event Process test"
  steps:
    - name: "testInput"
      type: "source"
      topic: "Test_topic_1"
      config: kafka
      destination: "transformation"
    - name: "transformation"
      type: "filter"
      modе: "all"
      filter:
      -  condition: "условие на DSL"
              destination: "destination1" # целевой шаг
          format: "json" # формат выходного сообщения
          log:
              level: "info"
              message: "Сообщение с id" || in.messageID || "отправлено в topic Topic_Name"
              input: environment.input_topic # опциональное поле
                          output: "topic1" # опциональное поле
        - condition: "условие на DSL"
                destination: "destination2" # целевой шаг
          format: "xml" # формат выходного сообщения
        - condition: "условие на DSL"
          destination: "decline" # при выполнении условия сообщение отбрасывается
      default:
          destination: "destinationDefault"
          format: "json"
    - name: "destination1"
      type: "destination"
      topic: "output1"
      config: kafka
    - name: "destination2"
      type: "destination"
      topic: "output2"
      config: kafka
    - name: "destinationDefault"
      type: "destination"
      topic: "output"
      config: 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"
  }
}