Структура файла потока «*.yaml»#

Конфигурацией потока трансформации является файл в формате YAML, путь до которого передается в аргументе запуска приложения --config.

Основным конфигурационным элементом является поле flow, которое содержит поля:

  1. name — имя задачи Flink, обязательное;

  2. steps — список шагов трансформации в потоке, каждый шаг должен ссылаться на следующий по имени из поля name каждого элемента списка в поле destination. Только шаги с типом destination не ссылаются на другие шаги из списка.

Пример файла *.yaml:

flow:
  name: "Event Process test"
  steps:
    - name: "testInput"
      type: "source"
      topic: "input"
      config:
        type: "kafka"
        "bootstrap.servers": "localhost:9092"
        consumer:
          "group.id": "event-process-flow-group"
          "client.id": "event-process-flow-client"
        producer:
          "client.id": "event-process-flow-producer"
      override:

      destination: "transformation"
    - name: "transformation"
      type: "dsl"
      format:
        input: "json"
        output: "json"
      dsl:
        path: "AsIs.tr"
      destination: "output"
    - name: "output"
      type: "destination"
      topic: "output"
      config:
        type: "kafka"
        "bootstrap.servers": "localhost:9092"
        consumer:
          "group.id": "event-process-flow-group"
          "client.id": "event-process-flow-client"
        producer:
          "client.id": "event-process-flow-producer"

Так же в файле *.yaml можно передать стендозависимые параметры через файл в аргументе запуска --defaults. В файле должно быть корневое поле defaults со значением объекта, каждое поле которого представляет собой имя пресета настроек подключения к транспорту, а значение представляет собой параметры подключения. Уже в шагах, где необходимо передать параметры подключения к транспорту или СУБД, можно передать вместо словаря параметров имя пресета из defaults. Так же в шагах source и destination появилось поле override, в котором можно переопределить отдельные параметры из defaults.

Пример файла *.yaml для параметра --defaults:

defaults:
  preset1:
    type: "kafka"
    "bootstrap.servers": "localhost:9092"
  preset2:
    type: "database"
    url: "jdbc:postgresql://localhost:5432/testDB"
    driver: "org.postgresql.Driver"

Пример файла *.yaml заполнения конфигурации с ссылками на defaults и перезапись параметром из него:

flow:
  name: "Event Process test"
  steps:
    - name: "testInput"
      type: "source"
      topic: "input"
      config: "preset1"
      override:
        "group.id": "event-process-flow-group"
        "client.id": "event-process-flow-client"
      destination: "transformation"
    - name: "transformation"
      type: "dsl"
      format:
        input: "json"
        output: "json"
      dsl:
        path: "AsIs.tr"
      database:
        config: "preset2"
        # Запросы, используемые в функциях с PreparedStatement
        statements:
          "selectTimeFromAccount": "SELECT time FROM accounts WHERE acctId = ?;"
          "selectAllFromAccounts": "SELECT * FROM accounts;"
          "deleteFromAccounts": "DELETE FROM accounts WHERE acctId = ?;"
          "mergeIntoAccounts": "MERGE INTO accounts VALUES (?, ?);"
        timeout: "100"
        retryCount: "5"
        retryInterval: "500"
        retryMultiplier: "2"
        typeCasting: true
      destination: "output"
    - name: "output"
      type: "destination"
      topic: "output"
      config: "preset1"
      override:
        "client.id": "event-process-flow-producer"