메트릭
메트릭 (Metrics)
PyFlink는 외부 시스템으로 메트릭을 수집하고 노출(expose)할 수 있는 메트릭 시스템을 제공해요.
출처: 문서
본문
메트릭 등록하기 (Registering metrics)
Python 사용자 정의 함수의 open 메서드에서 function_context.get_metric_group()을 호출해 메트릭 시스템에 접근할 수 있어요.
get_metric_group() 메서드는 MetricGroup 객체를 반환하며, 이 객체에 새 메트릭을 만들고 등록할 수 있어요.
메트릭 타입 (Metric types)
PyFlink는 Counter, Gauge, Distribution, Meter를 지원해요.
Counter
Counter는 무엇인가를 세는 데 사용돼요. 현재 값은 inc()/inc(n: int) 또는 dec()/dec(n: int)를 사용해 증가/감소시킬 수 있어요.
MetricGroup에서 counter(name: str)를 호출해 Counter를 만들고 등록할 수 있어요.
from pyflink.table.udf import ScalarFunction
class MyUDF(ScalarFunction):
def __init__(self):
self.counter = None
def open(self, function_context):
self.counter = function_context.get_metric_group().counter("my_counter")
def eval(self, i):
self.counter.inc(i)
return i
Gauge
Gauge는 요청 시(on demand) 값을 제공해요. MetricGroup에서 gauge(name: str, obj: Callable[[], int])를 호출해 gauge를 등록할 수 있어요. Callable 객체는 값을 보고하는 데 사용돼요. Gauge 메트릭은 정수 값만 허용해요.
from pyflink.table.udf import ScalarFunction
class MyUDF(ScalarFunction):
def __init__(self):
self.length = 0
def open(self, function_context):
function_context.get_metric_group().gauge("my_gauge", lambda : self.length)
def eval(self, i):
self.length = i
return i - 1
Distribution
보고된 값의 분포에 대한 정보(sum, count, min, max, mean)를 보고하는 메트릭이에요. 값은 update(n: int)를 사용해 갱신할 수 있어요. MetricGroup에서 distribution(name: str)을 호출해 distribution을 등록할 수 있어요. Distribution 메트릭은 정수 분포만 허용해요.
from pyflink.table.udf import ScalarFunction
class MyUDF(ScalarFunction):
def __init__(self):
self.distribution = None
def open(self, function_context):
self.distribution = function_context.get_metric_group().distribution("my_distribution")
def eval(self, i):
self.distribution.update(i)
return i - 1
Meter
Meter는 평균 처리량(throughput)을 측정해요. 이벤트의 발생은 mark_event() 메서드로 등록할 수 있어요. 동시에 여러 이벤트가 발생한 경우는 mark_event(n: int) 메서드로 등록할 수 있어요. MetricGroup에서 meter(self, name: str, time_span_in_seconds: int = 60)을 호출해 meter를 등록할 수 있어요.
time_span_in_seconds의 기본값은 60이에요.
from pyflink.table.udf import ScalarFunction
class MyUDF(ScalarFunction):
def __init__(self):
self.meter = None
def open(self, function_context):
# an average rate of events per second over 120s, default is 60s.
self.meter = function_context.get_metric_group().meter("my_meter", time_span_in_seconds=120)
def eval(self, i):
self.meter.mark_event(i)
return i - 1
스코프 (Scope)
Scope 정의에 대한 자세한 내용은 Java 메트릭 문서를 참고하세요.
사용자 스코프 (User Scope)
MetricGroup.add_group(key: str, value: str = None)을 호출해 사용자 스코프를 정의할 수 있어요.
value가 None이 아니면 새 키-값 MetricGroup 쌍이 만들어져요.
키 그룹은 이 그룹의 하위 그룹에 추가되고, 값 그룹은 키 그룹의 하위 그룹에 추가돼요. 이 경우 값 그룹이 반환되며 사용자 변수가 정의돼요.
function_context \
.get_metric_group() \
.add_group("my_metrics") \
.counter("my_counter")
function_context \
.get_metric_group() \
.add_group("my_metrics_key", "my_metrics_value") \
.counter("my_counter")
시스템 스코프 (System Scope)
시스템 스코프에 대한 자세한 내용은 Java 메트릭 문서를 참고하세요.
모든 변수 목록 (List of all Variables)
모든 변수 목록에 대한 자세한 내용은 Java 메트릭 문서를 참고하세요.
사용자 변수 (User Variables)
MetricGroup.addGroup(key: str, value: str = None)을 호출하고 value 매개변수를 지정해 사용자 변수를 정의할 수 있어요.
중요: 사용자 변수는 스코프 포맷에 사용할 수 없어요.
function_context \
.get_metric_group() \
.add_group("my_metrics_key", "my_metrics_value") \
.counter("my_counter")
PyFlink와 Flink 공통 부분 (Common part between PyFlink and Flink)
다음 섹션에 대한 자세한 내용은 Java 메트릭 문서를 참고하세요:
- 리포터 (Reporter).
- 시스템 메트릭 (System metrics).
- Latency 추적 (Latency tracking).
- REST API 통합 (REST API integration).
- 대시보드 통합 (Dashboard integration).