Структура файла потока «*.yaml»#
Конфигурацией потока трансформации является файл в формате YAML, путь до которого передается в аргументе запуска приложения --config.
Основным конфигурационным элементом является поле flow, которое содержит поля:
name— имя задачи Flink, обязательное;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"