트레이스
트레이스 (Traces)
Flink는 트레이스를 수집하고 외부 시스템으로 노출하는 추적(tracing) 시스템을 제공해요. 사용자 함수에서 Span을 보고하고, Flink가 보고하는 시스템 트레이스를 확인할 수 있어요.
출처: 문서
본문
Flink는 외부 시스템으로 트레이스를 수집·노출할 수 있는 추적 시스템을 제공해요.
트레이스 보고 (Reporting traces)
RichFunction을 확장하는 모든 사용자 함수에서 getRuntimeContext().getMetricGroup()을 호출해 추적 시스템에 접근할 수 있어요. 이 메서드는 MetricGroup 객체를 반환하며, 이를 통해 span 트리가 있는 새 단일 트레이스를 보고할 수 있어요.
단일 Span 보고 (Reporting single Span)
Span은 Flink에서 특정 시점에 특정 기간 동안 발생한 어떤 프로세스를 나타내며, TraceReporter에 보고돼요. Span을 보고하려면 MetricGroup#addSpan(SpanBuilder) 메서드를 사용할 수 있어요.
public class MyClass {
void doSomething() {
// (...)
metricGroup.addSpan(
Span.builder(MyClass.class, "SomeAction")
.setStartTsMillis(startTs) // Optional
.setEndTsMillis(endTs) // Optional
.setAttribute("foo", "bar") // Optional
.addChild(Span.builder(MyClass.class, "ChildAction") // Optional
.addChildren(List.of(
Span.builder(MyClass.class, "AnotherChildAction")); // Optional
}
}
Python에서 Span을 보고하는 것은 현재 지원되지 않아요.
Reporter
Flink의 trace reporter를 설정하는 방법에 대한 정보는 trace reporters 문서를 참고하세요.
시스템 트레이스 (System traces)
Flink는 아래 나열된 트레이스를 보고해요. 아래 표들은 일반적으로 5개의 컬럼을 가져요.
- "Scope" 컬럼은 해당 트레이스가 보고되는 범위가 무엇인지 설명해요.
- "Name" 컬럼은 보고된 트레이스의 이름을 설명해요.
- "Attributes" 컬럼은 주어진 트레이스와 함께 보고되는 모든 속성의 이름을 나열해요.
- "Description" 컬럼은 주어진 속성이 무엇을 보고하는지에 대한 정보를 제공해요.
체크포인팅 및 초기화 (Checkpointing and initialization)
Flink는 이벤트가 종료 상태인 COMPLETED 또는 FAILED에 도달하면 전체 체크포인트와 job 초기화 이벤트에 대해 단일 span 트레이스를 보고해요.
| Scope | Name | Attributes | Description |
|---|---|---|---|
| org.apache.flink.runtime.checkpoint.CheckpointStatsTracker | Checkpoint | startTs | 체크포인트가 시작된 타임스탬프 |
| endTs | 체크포인트가 종료된 타임스탬프 | ||
| checkpointId | 체크포인트의 Id | ||
| checkpointedSize | 이 체크포인트 동안 체크포인트된 상태의 크기(바이트). 증분 체크포인트를 사용하면 fullSize보다 작을 수 있어요. | ||
| fullSize | 이 체크포인트가 참조하는 상태의 전체 크기(바이트). 증분 체크포인트를 사용하면 checkpointSize보다 클 수 있어요. | ||
| checkpointStatus | 이 체크포인트의 상태: FAILED 또는 COMPLETED | ||
| checkpointType | 체크포인트의 유형. 예: "Checkpoint", "Full Checkpoint" 또는 "Terminate Savepoint" ... | ||
| isUnaligned | 체크포인트가 정렬되었는지 정렬되지 않았는지 여부 | ||
| JobInitialization | startTs | job 초기화가 시작된 타임스탬프 | |
| endTs | job 초기화가 종료된 타임스탬프 | ||
| checkpointId (optional) | job이 복구된 체크포인트의 Id (있는 경우) | ||
| fullSize | 복구 중 사용된 체크포인트가 참조하는 상태의 전체 크기(바이트, 있는 경우) | ||
| (Max/Sum)MailboxStartDurationMs | 모든 subtask에서 subtask가 생성된 시점부터 해당 subtask의 모든 클래스와 객체가 초기화될 때까지의 기간(최대·합계) 집계 | ||
| (Max/Sum)ReadOutputDataDurationMs | 모든 subtask에서 정렬되지 않은 체크포인트의 출력 버퍼를 읽는 기간(최대·합계) 집계 | ||
| (Max/Sum)InitializeStateDurationMs | 모든 subtask에서 상태 백엔드를 초기화하는 기간(최대·합계) 집계 (상태 파일 다운로드 시간 포함) | ||
| (Max/Sum)GateRestoreDurationMs | 모든 subtask에서 정렬되지 않은 체크포인트의 입력 버퍼를 읽는 기간(최대·합계) 집계 | ||
| (Max/Sum)DownloadStateDurationMs(optional - currently only supported by RocksDB Incremental) | 모든 subtask에서 DFS에서 상태 파일을 다운로드하는 기간(최대·합계) 집계 | ||
| (Max/Sum)RestoreStateDurationMs(optional - currently only supported by RocksDB Incremental) | 모든 subtask에서 완전히 로컬화된 상태, 즉 모든 원격 상태가 다운로드된 후 상태 백엔드를 복원하는 기간(최대·합계) 집계 | ||
| (Max/Sum)RestoredStateSizeBytes.[location] | 모든 subtask에서 위치별로 복원된 상태 크기(최대·합계) 집계. 가능한 위치는 Enum StateObjectSizeStatsCollector에 LOCAL_MEMORY, LOCAL_DISK, REMOTE, UNKNOWN으로 정의돼 있어요. |
||
| (Max/Sum)RestoreAsyncCompactionDurationMs(optional - currently only supported by RocksDB Incremental) | 증분 복원 후 비동기 압축을 위한 모든 subtask의 기간(최대·합계) 집계 |