Watermark
Watermark
참고: DataStream API V2는 기존 DataStream API를 점진적으로 대체할 새로운 API 집합입니다. 현재 실험 단계이며 프로덕션에는 완전히 사용할 수 없습니다.
Watermark를 소개하기 전에, DataStream V2의 Watermark는 이벤트 시간의 진행을 측정하는 기존 Watermark를 의미하지 않는다는 점을 알아야 합니다. V2의 Watermark는 사용자가 사용자 지정할 수 있고 스트림을 따라 전파될 수 있는 특별한 이벤트입니다.
Flink에서 Watermark를 사용하는 것은 세 가지 핵심 단계를 포함합니다.
- Watermark 정의 및 선언
- Watermark 내보내기
- Watermark 처리
이 단계들을 따라 Flink에서 Watermark를 사용하는 방법을 알아봅니다.
출처: 문서
본문
Watermark 정의 및 선언
워터마크는 데이터를 운반하는 특별한 이벤트입니다. Source 또는 ProcessFunction에서 내보내질 수 있고, 스트림을 따라 전파되며, 다운스트림 ProcessFunction이 수신합니다. Watermark를 정의할 때 고려할 네 가지 측면이 있습니다.
-
[필수] Watermark 식별자(Identifier)
스트림에 여러 유형의 Watermark가 전파될 수 있으므로 각 Watermark를 구분하기 위해 식별자를 할당하는 것이 중요합니다.
식별자는 String이어야 하고 대소문자를 구분하며 전체 작업 내에서 전역적으로 고유해야 합니다.
-
[필수] Watermark 데이터 타입
Watermark의 데이터 타입을 지정하는 것이 중요합니다. 현재 Flink는 Long과 Bool 두 가지 타입을 지원합니다.
-
[필수] 결합 함수(Combine Function)와
combineWaitForAllChannelsProcessFunction은 여러 업스트림 입력이 있을 수 있고 각각 병렬도가 다를 수 있으므로 서로 다른 입력 채널에서 여러 Watermark를 받을 수 있습니다. 이러한 경우 사용자는ProcessFunction으로 출력하기 전에 입력 채널의 Watermark를 결합하려는 경우가 많습니다.Flink는 다음 결합 함수를 지원합니다.
Long타입 Watermark:MIN: 수신된 모든 워터마크의 최소값을 유지하고 출력합니다.MAX: 수신된 모든 워터마크의 최대값을 유지하고 출력합니다.
Bool타입 Watermark:AND: 수신된 모든 워터마크의 논리 AND 결과를 유지하고 출력합니다.OR: 수신된 모든 워터마크의 논리 OR 결과를 유지하고 출력합니다.
또한 사용자는 결합 과정이
ProcessFunction이 모든 업스트림 채널에서 Watermark를 받을 때까지 기다릴지 여부를 구성할 수 있습니다. 이는 일부 시나리오에서 특히 유용합니다. 예를 들어 이벤트 시간 워터마크는 모든 입력에서 워터마크를 받을 때까지 기다린 다음 결합해야 합니다. 이렇게 하면 이벤트 시간 워터마크가 운반하는 시간이 감소하지 않도록 보장됩니다. 기본적으로combineWaitForAllChannels설정은 false입니다. -
[선택] 프레임워크의
WatermarkHandlingStrategyWatermarkHandlingStrategy는 사용자 정의ProcessFunction이 Watermark를 pop하지 않을 때 프레임워크가 워터마크를 다운스트림ProcessFunction으로 보내야 하는지 결정합니다. 이 설정에는 두 가지 옵션이 있습니다.- IGNORE: 프레임워크가 어떤 조치도 취하지 않아야 합니다.
- FORWARD: 프레임워크가 워터마크를 다운스트림으로 보내야 합니다.
이 선택적 설정은 일부 경우에 유용합니다. 예를 들어
IGNORE로 설정하면 프레임워크가 이 Watermark를 전파할 필요가 없으며, 전송 제어는 사용자에게 달려 있음을 나타낼 수 있습니다.
Watermark 정의 과정을 단순화하기 위해 Flink는 WatermarkBuilder를 제공합니다. 이 빌더는 궁극적으로 WatermarkDeclaration 객체를 생성합니다. 다음은 이를 사용해 Watermark를 정의하는 예시입니다.
LongWatermarkDeclaration watermarkDeclaration = WatermarkDeclarations
.newBuilder("MY_CUSTOM_WATERMARK_IDENTIFIER")
.typeLong()
.combineFunctionMax()
.combineWaitForAllChannels(true)
.defaultHandlingStrategyForward()
.build();
사용자가 Watermark를 정의했다면 ProcessFunction#declareWatermarks 또는 Source#declareWatermarks에서 선언하는 것이 중요합니다. 이 단계를 통해 프레임워크가 이를 제대로 인식할 수 있습니다. ProcessFunction에서 Watermark를 선언하는 예시입니다.
public class CustomProcessFunction
implements OneInputStreamProcessFunction<Long, Long> {
LongWatermarkDeclaration watermarkDeclaration = ...;
@Override
public Set<? extends WatermarkDeclaration> declareWatermarks() {
return Set.of(watermarkDeclaration);
}
}
각 유형의 Watermark는 작업에서 한 번만 선언하면 됩니다.
Watermark 내보내기
사용자는 Watermark 선언을 사용해 워터마크를 생성할 수 있습니다. 값이 1인 Long 타입 워터마크를 생성하는 예시입니다.
LongWatermarkDeclaration watermarkDeclaration = ...;
LongWatermark watermark = watermarkDeclaration.newWatermark(1);
사용자는 nonPartitionedContext.getWatermarkManager().emitWatermark(watermark)를 호출해 ProcessFunction에서 Watermark를 내보낼 수 있고, sourceReaderContext.emitWatermark(watermark)를 호출해 Source에서도 Watermark를 내보낼 수 있습니다. 다음은 ProcessFunction에서 Watermark를 내보내는 예시입니다.
public class CustomProcessFunction
implements OneInputStreamProcessFunction<Long, Long> {
LongWatermarkDeclaration watermarkDeclaration = ...;
@Override
public Set<? extends WatermarkDeclaration> declareWatermarks() {
return Set.of(watermarkDeclaration);
}
@Override
public void processRecord(Long record, Collector<Long> output, PartitionedContext<Long> ctx)
throws Exception {
// do something as needed
// generate and emit Watermark
LongWatermark watermark = watermarkDeclaration.newWatermark(1);
ctx.getNonPartitionedContext().getWatermarkManager().emitWatermark(watermark);
}
}
Watermark 처리
ProcessFunction이 Watermark를 수신하면 프레임워크는 ProcessFunction#onWatermark 메서드를 호출해 처리합니다. 따라서 사용자는 Watermark를 적절히 처리하기 위해 ProcessFunction#onWatermark를 재정의해야 합니다.
ProcessFunction#onWatermark 메서드의 반환 값은 WatermarkHandlingResult 타입이며 두 가지 옵션이 있습니다.
- PEEK:
ProcessFunction이 Watermark를 peek만 하고, 이 워터마크를 처리하는 것은 프레임워크의 책임입니다.- 프레임워크는 워터마크와 연관된
WatermarkHandlingStrategy에 따라 워터마크를 전달하거나 무시합니다.
- 프레임워크는 워터마크와 연관된
- POLL: 이 Watermark는 process function 자체가 다운스트림으로 보내야 합니다. 프레임워크는 추가 처리를 하지 않습니다.
기본적으로 ProcessFunction#onWatermark 메서드는 WatermarkHandlingResult.PEEK를 반환합니다.
ProcessFunction에서 Watermark를 처리하는 예시입니다.
public static final String CUSTOM_WATERMARK_IDENTIFIER = "CUSTOM_WATERMARK_IDENTIFIER";
public class CustomProcessFunction
implements OneInputStreamProcessFunction<Long, Long> {
...
@Override
public WatermarkHandlingResult onWatermark(
Watermark watermark, Collector<Long> output, NonPartitionedContext<Long> ctx) {
// For Watermark that this ProcessFunction is interested in, process the watermark
if (watermark.getIdentifier().equals(CUSTOM_WATERMARK_IDENTIFIER)) {
// do something as needed
}
// For other Watermarks, return PEEK
return WatermarkHandlingResult.PEEK;
}
}