Обработка ошибок в шагах трансформации и агрегации событий#

Во время преобразования форматов событий при возникновении ошибки обработки события может осуществляться перехват ошибок. Это происходит только в процессе трансформации накопленных событий в окне агрегации. При этом ошибки при вычислении ключа и выполнении триггера приводят к ошибке обработки. На этапе вычисления ключа агрегации и выполнения триггера перехват не выполняется из-за ограничений API Apache Flink. Данное поведение для обработчика является штатным.

В настройках шагов трансформации и агрегации событий можно задать поле error, в котором необходимо указать первый шаг побочной ветки для дальнейшей обработки или пересылки.

При передаче событий в эту ветку в заголовки сообщения Kafka будут добавлены следующие записи:

  1. stepName имя шага, на котором произошла ошибка обработки события.

  2. errorType тип ошибки, принимает следующие значения:

    • user ошибка выполнения правил трансформации, в том числе вызов функции exception();

    • system системная ошибка обработчика, например, ошибка чтения форматов JSON, XML, CSV или AVRO.

  3. errorMessage сообщение об ошибке.

Пример:

{
  source: {
    name: "testInput"
    type: "test"
    topic: "input"
    config: ${defaults.kafka}
    destination: ${aggregationStep}
  }

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

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

  aggregationStep: {
    name: "dsl-aggregation"
    type: "aggregation"
    format: {
      input: "json"
      output: "json"
    }
    key {
      type: "dsl"
      path: "key.tr"
    }
    trigger: {
      type: "dsl"
      path: "trigger.tr"
      timeout: 100
    }
    dsl: {
      path: "simple.tr"
    }
    destination: ${destination}
    error: ${deadletter}
  }

  flow: {
    name: "Event Process test"
    source: [${source}]
  }
}