메트릭

메트릭 (Metrics)

Flink는 외부 시스템으로 메트릭을 수집하고 노출(expose)할 수 있는 메트릭 시스템을 제공해요.

출처: 문서

본문

이 페이지에서는 메트릭 등록, 스코프, 리포터, 시스템 메트릭, 지연 추적, REST API 통합과 대시보드 통합에 대해 다뤄요. 자세한 코드 예시는 아래에 분포해 있어요.

메트릭 등록 (Registering metrics)

RichFunction을 상속하는 어떤 사용자 함수에서든 getRuntimeContext().getMetricGroup()을 호출해 메트릭 시스템에 접근할 수 있어요. 이 메서드는 새 메트릭을 만들고 등록할 수 있는 MetricGroup 객체를 반환해요.

Flink는 Counters, Gauges, Histograms, Meters를 지원해요. 각각에 대해 Java와 Python 예시가 제공돼요.

Counter

Counter는 무엇인가를 세는 데 사용돼요. 현재 값은 inc()/inc(long n) 또는 dec()/dec(long n)을 사용해 증가/감소시킬 수 있어요. MetricGroup에서 counter(String name)을 호출해 Counter를 만들고 등록할 수 있어요.

public class MyMapper extends RichMapFunction<String, String> {
  private transient Counter counter;
  @Override
  public void open(OpenContext ctx) {
    this.counter = getRuntimeContext().getMetricGroup().counter("myCounter");
  }
  @Override
  public String map(String value) throws Exception {
    this.counter.inc();
    return value;
  }
}

또한 counter(String name, Counter counter)로 자신만의 Counter 구현을 사용할 수도 있어요.

Gauge

Gauge는 어떤 타입의 값이라도 요청 시 제공해요. org.apache.flink.metrics.Gauge 인터페이스를 구현하는 클래스를 만든 후 gauge(String name, Gauge gauge)으로 등록해요. 반환되는 값의 타입에는 제한이 없어요.

Histogram

Histogram은 long 값의 분포를 측정해요. histogram(String name, Histogram histogram)으로 등록할 수 있어요. Flink는 Histogram 기본 구현을 제공하지 않지만, Codahale/DropWizard 히스토그램을 사용할 수 있는 래퍼 flink-metrics-dropwizard를 제공해요.

Meter

Meter는 평균 처리량을 측정해요. markEvent()로 이벤트 발생을 등록하고 meter(String name, Meter meter)으로 등록할 수 있어요.

스코프 (Scope)

모든 메트릭에는 식별자와 키-값 쌍 집합이 할당돼요. 식별자는 장명 등록 시 사용자 정의 이름, 선택적 사용자 정의 스코프, 시스템 제공 스코프의 세 컴포넌트를 기반으로 해요. 예를 들어 A.B가 시스템 스코프, C.D가 사용자 스코프, E가 이름이면 식별자는 A.B.C.D.E예요. metrics.scope.delimiter 키로 구분자를 구성할 수 있어요(기본: .).

사용자 스코프는 MetricGroup#addGroup 메서드로, 시스템 스코프는 metrics.scope.* 구성 키로 정의해요. metrics.scope.operator 기본값은 .taskmanager.<job_name>.<operator_name>.<subtask_index> 형태로, localhost.taskmanager.1234.MyJob.MyOperator.0.MyMetric 같은 식별자를 만들어요.

사용 가능한 변수: JobManager (<hostname>, <jm_id>), TaskManager (<hostname>, <tm_id>), Job (<job_id>, <job_name>), Task (<task_id>, <task_name>, <task_attempt_num>, <subtask_index>), Operator (<operator_name>, <subtask_index>).

사용자 변수는 MetricGroup#addGroup(String key, String value)로 정의할 수 있지만, 스코프 포맷에는 사용할 수 없어요. Transformation.addMetricVariable을 사용해 연산자별 커스텀 변수를 정의할 수도 있어요.

리포터 (Reporter)

리포터 설정은 메트릭 리포터 문서를 참조하세요.

시스템 메트릭 (System metrics)

기본적으로 Flink는 현재 상태에 대한 깊은 통찰을 제공하는 여러 메트릭을 수집해요. CPU, 메모리, 스레드, 클래스로더, 체크포인트, 상태, 네트워크, 연산자 등 다양한 범주가 있어요. 예를 들어 Operator 범주에는 비동기 상태 처리 메트릭이 있어요:

Scope Infix Metrics Description Type
Operator asyncStateProcessing numInFlightRecords 비동기 실행 컨트롤러 버퍼의 진행 중(in-flight) 레코드 수. Gauge
activeBufferSize 처리 대기 중인 레코드 수. Gauge
blockingBufferSize 진행 중인 레코드에 의해 차단된 레코드 수. Gauge
numBlockingKeys 비동기 실행 컨트롤러에서 차단된 서로 다른 키 수. Gauge

종단 간 지연 추적 (End-to-End latency tracking)

Flink는 시스템을 통과하는 레코드의 지연을 추적할 수 있게 해줘요. 기본적으로 비활성화돼 있으며, latencyTrackingInterval을 양수로 설정해 활성화해요. 소스는 주기적으로 LatencyMarker를 방출하고, 이를 사용해 소스와 각 다운스트림 연산자 간의 지연 분포(히스토그램 메트릭)를 파생해요. 이 기능은 클러스터 성능에 크게 영향을 줄 수 있으므로 디버깅 목적으로만 사용하는 것을 권장해요.

상태 접근 지연 추적 (State access latency tracking)

state.latency-track.keyed-state-enabled을 true로 설정해 키처리 상태 접근 지연 추적을 활성화할 수 있어요. state.latency-track.sample-interval(기본 100)과 state.latency-track.history-size(기본 128)로 샘플링과 히스토리 크기를 제어해요.

상태 키/값 크기 추적 (State key/value size tracking)

state.size-track.keyed-state-enabled을 true로 설정해 키/값 크기 추적을 활성화할 수 있어요. state.size-track.sample-intervalstate.size-track.history-size로 제어해요.

REST API 통합 (REST API integration)

메트릭은 Monitoring REST API로 질의할 수 있어요. 예: /jobmanager/metrics, /taskmanagers/<id>/metrics, /jobs/<jobid>/metrics 등. 특수 문자(예: #, $, &, +, /, ;, =, @)는 URL 인코딩으로 이스케이프해야 해요.

대시보드 통합 (Dashboard integration)

각 태스크/연산자에 대해 수집된 메트릭은 대시보드에서도 시각화할 수 있어요. 작업 페이지에서 Metrics 탭을 선택하고 Add Metric 메뉴로 표시할 메트릭을 선택해요. 각 메트릭은 별도 그래프로 시각화되며 10초마다 자동 갱신돼요. 숫자 메트릭만 시각화할 수 있어요.

더 알아보기 (Learn more)