Простой поток агрегации#

Simple Aggregation.png

Примеры агрегации приведены в разделах Поток со слиянием, ветвлением и агрегацией и Поток с ветвлением, агрегацией и слиянием.

Особенность использования агрегации: если в блоке агрегации был задан тайм-аут, то по его окончании будет выполнена обработка всех накопленных событий независимо от срабатывания триггера. Данную ситуацию необходимо учитывать при написании логики обработки результатов если используется тайм-аут.

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

  source: {
    name: "simple aggregation"
    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: ${aggregationStep}
  }

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

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