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

Merge Branch Aggregation.png

{

  flow: {
    name: "JobName"
    source: [${sourceRQ},${sourceRS}] # Массив имен блоков-источников
  }

# Первый источник
  sourceRQ: {
    name: "source request"
    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}
  }

# Второй источник
  sourceRS: {
    name: "source response"
    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}
  }

# Блок слияния потоков источников
  mergeStep: {
    name: "merge sources"
    type: "merge"
    destination: ${transformStep}
  }

# Блок трансформации с ветвлением. Выдает на выход одно из выходных сообщений. Также направляет сообщения на агрегацию. Это разделение называется бранчированием
  transformStep: {
    name: "transformation step"
    type: "dsl"
    dsl: {
      source: "file"
      path: "event.tr" # DSL-файл, описывает трансформацию
    }
    format : {
      input: "xml"
      output: "xml"
    }
    destination: {
      type: "branch"
      destination: ${destinationRQ} # Основная ветвь потока
      branches: [${aggregationStep}] # Массив побочных ветвей потока
    }
  }

# Блок агрегации
  aggregationStep: {
    name: "aggregation step" # По этому имени в DSL блока трансформации указывается отправка события в ветвь (branch) агрегации
    type: "aggregation"
    format: {
      input: "xml"
      output: "xml"
    }
    key {
      type: "dsl"
      path: "eventKey.tr"
    }
    trigger: {
      type: "dsl"
      path: "eventTrigger.tr"
      timeout: 3600000
    }
    dsl: {
      path: "aggregationEvent.tr"
    }
    destination: ${destinationRS}
  }

# Блок-получатель одного из выходных событий (после трансформации)
  destinationRQ: {
    name: "destination request" # По этому имени в DSL блока трансформации указывается отправка события в выходную ветвь (branch)
    type: "destination"
    topic: "output request"
    config: ${defaults.kafka}
  }

# Блок-получатель другого типа выходных событий (после агрегации)
  destinationRS: {
    name: "destination response"
    type: "destination"
    topic: "output response"
    config: ${defaults.kafka}
  }
}