Job 상태 변경 리스너

Job 상태 변경 리스너 (Job status changed listener)

Flink는 사용자가 작업 상태 변경을 처리하는 사용자 정의 로직을 등록할 수 있는 플러그형(pluggable) 인터페이스를 제공합니다. 이때 소스/싱크에 대한 계보(lineage) 정보도 제공됩니다. 이를 통해 사용자는 자신만의 Flink lineage reporter를 구현해 Datahub나 Openlineage 같은 서드파티 데이터 계보 시스템에 계보 정보를 보낼 수 있습니다.

출처: 문서

본문

Flink는 작업 상태 변경을 처리하기 위한 플러그형 인터페이스를 제공하며, 여기에는 소스/싱크에 대한 계보 정보가 포함됩니다. 이를 통해 사용자는 Datahub나 Openlineage 같은 서드파티 데이터 계보 시스템에 계보 정보를 보내는 자신만의 Flink lineage reporter를 구현할 수 있습니다.

Job status changed listener는 애플리케이션의 상태가 변경될 때마다 트리거됩니다. 데이터 계보 정보는 JobCreatedEvent에 포함되어 있습니다.

Job status changed listener 플러그인 구현하기

사용자 정의 JobStatusChangedListener 플러그인을 구현하려면 다음을 수행해야 합니다.

  • JobStatusChangedListener 인터페이스를 구현해 자신만의 JobStatusChangedListener를 추가합니다.
  • JobStatusChangedListenerFactory 인터페이스를 구현해 자신만의 JobStatusChangedListenerFactory를 추가합니다.
  • 서비스 항목을 추가합니다. META-INF/services/org.apache.flink.core.execution.JobStatusChangedListenerFactory 파일을 만들고, 이 파일 안에 job status changed listener 팩토리 클래스의 클래스 이름을 넣습니다(자세한 내용은 Java Service Loader 문서 참조).

그런 다음 JobStatusChangedListener, JobStatusChangedListenerFactory, META-INF/services/와 모든 외부 의존성을 포함하는 jar을 만듭니다. Flink 배포의 plugins/ 안에 임의의 이름으로 디렉터리를 만들고(예: "job-status-changed-listener") 그 jar을 이 디렉터리에 넣습니다. 자세한 내용은 Flink Plugin을 참조하세요.

JobStatusChangedListenerFactory 예시:

package org.apache.flink.test.execution;

public static class TestingJobStatusChangedListenerFactory
        implements JobStatusChangedListenerFactory {

    @Override
    public JobStatusChangedListener createListener(Context context) {
        return new TestingJobStatusChangedListener();
    }
}

JobStatusChangedListener 예시:

package org.apache.flink.test.execution;

private static class TestingJobStatusChangedListener implements JobStatusChangedListener {

    @Override
    public void onEvent(JobStatusChangedEvent event) {
        statusChangedEvents.add(event);
    }
}

구성 (Configuration)

Flink 컴포넌트는 시작 시 JobStatusChangedListener 플러그인을 로드합니다. 여러분의 JobStatusChangedListener가 로드되도록 하려면 모든 클래스 이름을 execution.job-status-changed-listeners의 일부로 정의해야 합니다. 이 구성이 비어 있으면 어떤 enricher도 시작되지 않습니다. 예시:

    execution.job-status-changed-listeners = org.apache.flink.test.execution.TestingJobStatusChangedListenerFactory

더 알아보기 (Learn more)