Поток с ветвлением, агрегацией и слиянием#

Branch Aggregation Merge.png

{
  flow: {
    name: "JobName"
    source: [${sourceST}] # Если источник один - массив указывать необязательно. Можно было написать тут просто "${sourceST}"
  }

  sourceST: {
    name: "simple transfer"
    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: "dsl"
    dsl: {
      source: "file"
      path: "event.tr"
    }
    format : {
      input: "xml"
      output: "xml"
    }
    destination: {
      type: "branch"
      destination: ${destination} # Основная ветвь потока
      branches: [${mergeStep}] # Массив побочных ветвей
    }
  }

# Блок агрегации
  aggregationStep: {
    name: "aggregation step" # По этому имени в DSL блока трансформации указывается отправка события в ветвь (branch) агрегации
    type: "aggregation"
    format: {
      input: "xml"
      output: "xml"
    }
    key: {
      type: "dsl"
      path: "eventKey.tr" # DSL-файл ключа агрегирования. Возвращает значение, по которому связываются события
    }
    trigger: {
      type: "dsl"
      path: "eventTrigger.tr" # DSL-файл, триггер агрегации. Описывает условие, по которому начинается обработка накопленных событий
      timeout: 10800000 # Время жизни "окна" агрегации. Считается от последнего пришедшего в блок сообщения
    }
    dsl: {
      path: "aggregationEvent.tr" # Основной DSL-файл агрегации. Описывает обработку накопленных событий
    }
    destination: ${mergeStep}
  }

# Блок слияния
  mergeStep: {
    name: "merge request and response"
    type: "merge"
    destination: ${destination}
  }

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