Создание и объединение нескольких веток обработки событий#

Порождение нескольких веток обработки событий#

Шаг предназначен для разделения и обработки событий по различным шагам в зависимости от заданных условий.

Создать побочные ветки обработки можно с помощью объекта с типом branch со следующими полями:

  1. destination — следующий шаг основной ветви обработки событий

  2. branches — массив следующих шагов в побочных ветках

destination: {
      type: "branch"
      destination: ${mainNextStep}
      branches: [${nextStepInFirstBranch}, ${nextStepInSecondBranch}]
    }

Создание побочных веток обработки#

Для маршрутизации событий в заголовок BranchName в файле .tr должно быть установлено имя ветки из массива branches файла .conf.

Значения в массиве branches файла .conf являются значениями поля name из последующих шагов.

Пример файла .tr:

headers[>].BranchName = "transformSecond"
out[>].first = in
headers[+>].BranchName = "transformThird"
out[+>].second = in

Пример файла .conf:

{
  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"
    type: "dsl"
    format: {
      input: "json"
      output: "json"
    }
    dsl: {
      path: "transform.tr"
    }
    destination: {
      type: "branch"
      destination: ${destination}
      branches: [${transformSecond}, ${transformThird}] // Значения в массиве `branches` должно быть установлено в заголовок `BranchName` из файла **.tr**. Также значения в массиве `branches` являются значениями поля `name` из последующих шагов `transformSecondStep` и `transformThirdStep`
    }
  }

  transformSecondStep: {
    name: "transformSecond" // Значение из массива `branches` блока трансформации `transformStep`
    type: "dsl"
    format: {
      input: "json"
      output: "json"
    }
    dsl: {
      path: "second.tr"
    }
    destination: ${destination1}
  }

  transformThirdStep: {
    name: "transformThird" // Значение из массива `branches` блока трансформации `transformStep`
    type: "dsl"
    format: {
      input: "json"
      output: "json"
    }
    dsl: {
      path: "third.tr"
    }
    destination: ${destination2}
  }

  destination1: {
    name: "destination1"
    type: "destination"
    topic: "output1"
    config: ${defaults.kafka}
  }

  destination2: {
    name: "destination2"
    type: "destination"
    topic: "output2"
    config: ${defaults.kafka}
  }
}

Маршрутизация потока по условию#

Существует возможность маршрутизиции событий в различные ветки по условию без заполнения файла .tr. Описано в разделе Трансформация событий в подразделе «Шаг фильтрации/маршрутизации».

Объединение нескольких веток обработки#

Шаг предназначен для объединения веток обработки событий в один поток. Все шаги, содержащие этот шаг в поле destination, будут объединены в один поток.

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

  1. name — имя шага

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

Объединение источников событий#

mergeStep: {
    name: "merge sources"
    type: "merge"
    destination: {}
  }

Пример файла .conf:

{

  flow: {
    name: "JobName"
    source: ${source}
  }

  sourceFirst: {
    name: "sourceFirst"
    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: ${mergeStep}
  }

  sourceSecond: {
    name: "sourceFirst"
    type: "source"
    topic: "input_second"
    config: {
      consumer: {
        "group.id": "event-process-flow-group"
        "client.id": "event-process-flow-client"
      }
      producer: {
        "client.id": "event-process-flow-producer"
      }
    }
    destination: ${mergeStep}
  }

  mergeStep: {
    name: "merge"
    type: "merge"
    destination: ${transformStep}
  }

  transformStep: {
    name: "transformation"
    type: "dsl"
    format: {
      input: "json"
      output: "json"
    }
    dsl: {
      path: "AsIs.tr"
    }
    destination: ${destination}
  }

  destination: {
    name: "destination"
    type: "destination"
    topic: "output1"
    config: ${defaults.kafka}
  }
}