커스텀 실패 인리처

커스텀 실패 인리처 (Custom Failure Enrichers)

Flink는 사용자가 커스텀 로직을 등록하고 실패(failure)에 추가 메타데이터 레이블(문자열 키-값 쌍)을 붙여 풍부하게 만들 수 있는 플러그형(pluggable) 인터페이스를 제공해요. 이를 통해 사용자는 자신만의 실패 인리치먼트(failure enrichment) 플러그인을 구현해 작업(job) 실패를 분류하고, 커스텀 메트릭을 노출하거나, 외부 알림 시스템을 호출할 수 있어요.

출처: 문서

본문

FailureEnricher는 JobManager에서 런타임에 예외가 보고될 때마다 트리거돼요. 각 FailureEnricher는 실패와 연관된 레이블을 비동기적으로 반환할 수 있으며, 이 레이블들은 JobManager의 REST API를 통해 노출돼요(예: type:System 레이블은 해당 실패가 시스템 오류로 분류됨을 의미).

커스텀 인리처 플러그인 구현하기

커스텀 FailureEnricher 플러그인을 구현하려면 다음을 수행해야 해요:

  • FailureEnricher 인터페이스를 구현해 자신만의 FailureEnricher를 추가해요.
  • FailureEnricherFactory 인터페이스를 구현해 자신만의 FailureEnricherFactory를 추가해요.
  • 서비스 엔트리를 추가해요. META-INF/services/org.apache.flink.core.failure.FailureEnricherFactory 파일을 만들고 그 안에 자신의 실패 인리처 팩토리 클래스 이름을 넣어요(자세한 내용은 Java Service Loader 문서 참고).

그런 다음 자신의 FailureEnricher, FailureEnricherFactory, META-INF/services/와 모든 외부 의존성을 포함하는 jar을 만들어요. Flink 배포판의 plugins/ 안에 임의의 이름(예: "failure-enrichment")으로 디렉터리를 만들고 그 jar를 이 디렉터리에 넣어요. 자세한 내용은 Flink Plugin 문서를 참고하세요.

모든 FailureEnricher는 값과 연관될 수 있는 출력 키(output keys) 집합을 정의해야 한다는 점에 주의하세요. 이 키 집합은 고유해야 하며, 그렇지 않으면 키가 겹치는 모든 인리처는 무시돼요.

FailureEnricherFactory 예시:

package org.apache.flink.test.plugin.jar.failure;

public class TestFailureEnricherFactory implements FailureEnricherFactory {

   @Override
   public FailureEnricher createFailureEnricher(Configuration conf) {
        return new CustomEnricher();
   }
}

FailureEnricher 예시:

package org.apache.flink.test.plugin.jar.failure;

public class CustomEnricher implements FailureEnricher {
    private final Set<String> outputKeys;
    
    public CustomEnricher() {
        this.outputKeys = Collections.singleton("labelKey");
    }

    @Override
    public Set<String> getOutputKeys() {
        return outputKeys;
    }

    @Override
    public CompletableFuture<Map<String, String>> processFailure(
            Throwable cause, Context context) {
        return CompletableFuture.completedFuture(Collections.singletonMap("labelKey", "labelValue"));
    }
}

설정 (Configuration)

JobManager는 시작 시 FailureEnricher 플러그인을 로드해요. 자신의 FailureEnricher가 로드되게 하려면 모든 클래스 이름을 jobmanager.failure-enrichers 설정의 일부로 정의해야 해요. 이 설정이 비어 있으면, 어떠한 인리처도 시작되지 않아요. 예시:

    jobmanager.failure-enrichers = org.apache.flink.test.plugin.jar.failure.CustomEnricher

검증 (Validation)

자신의 FailureEnricher가 로드되었는지 검증하려면 JobManager 로그에서 아래 줄을 확인할 수 있어요:

    Found failure enricher org.apache.flink.test.plugin.jar.failure.CustomEnricher at jar:file:/path/to/flink/plugins/failure-enrichment/flink-test-plugin.jar!/org/apache/flink/test/plugin/jar/failure/CustomEnricher.class

또한 JobManager의 REST API에서 failureLabels 필드를 찾아 질의할 수도 있어요:

    "failureLabels": {
        "labelKey": "labelValue"
    }

더 알아보기 (Learn more)