Watermark

Watermark

참고: DataStream API V2는 기존 DataStream API를 점진적으로 대체할 새로운 API 집합입니다. 현재 실험 단계이며 프로덕션에는 완전히 사용할 수 없습니다.

Watermark를 소개하기 전에, DataStream V2의 Watermark는 이벤트 시간의 진행을 측정하는 기존 Watermark를 의미하지 않는다는 점을 알아야 합니다. V2의 Watermark는 사용자가 사용자 지정할 수 있고 스트림을 따라 전파될 수 있는 특별한 이벤트입니다.

Flink에서 Watermark를 사용하는 것은 세 가지 핵심 단계를 포함합니다.

  1. Watermark 정의 및 선언
  2. Watermark 내보내기
  3. Watermark 처리

이 단계들을 따라 Flink에서 Watermark를 사용하는 방법을 알아봅니다.

출처: 문서

본문

Watermark 정의 및 선언

워터마크는 데이터를 운반하는 특별한 이벤트입니다. Source 또는 ProcessFunction에서 내보내질 수 있고, 스트림을 따라 전파되며, 다운스트림 ProcessFunction이 수신합니다. Watermark를 정의할 때 고려할 네 가지 측면이 있습니다.

  1. [필수] Watermark 식별자(Identifier)

    스트림에 여러 유형의 Watermark가 전파될 수 있으므로 각 Watermark를 구분하기 위해 식별자를 할당하는 것이 중요합니다.

    식별자는 String이어야 하고 대소문자를 구분하며 전체 작업 내에서 전역적으로 고유해야 합니다.

  2. [필수] Watermark 데이터 타입

    Watermark의 데이터 타입을 지정하는 것이 중요합니다. 현재 Flink는 Long과 Bool 두 가지 타입을 지원합니다.

  3. [필수] 결합 함수(Combine Function)와 combineWaitForAllChannels

    ProcessFunction은 여러 업스트림 입력이 있을 수 있고 각각 병렬도가 다를 수 있으므로 서로 다른 입력 채널에서 여러 Watermark를 받을 수 있습니다. 이러한 경우 사용자는 ProcessFunction으로 출력하기 전에 입력 채널의 Watermark를 결합하려는 경우가 많습니다.

    Flink는 다음 결합 함수를 지원합니다.

    • Long 타입 Watermark:
      • MIN: 수신된 모든 워터마크의 최소값을 유지하고 출력합니다.
      • MAX: 수신된 모든 워터마크의 최대값을 유지하고 출력합니다.
    • Bool 타입 Watermark:
      • AND: 수신된 모든 워터마크의 논리 AND 결과를 유지하고 출력합니다.
      • OR: 수신된 모든 워터마크의 논리 OR 결과를 유지하고 출력합니다.

    또한 사용자는 결합 과정이 ProcessFunction이 모든 업스트림 채널에서 Watermark를 받을 때까지 기다릴지 여부를 구성할 수 있습니다. 이는 일부 시나리오에서 특히 유용합니다. 예를 들어 이벤트 시간 워터마크는 모든 입력에서 워터마크를 받을 때까지 기다린 다음 결합해야 합니다. 이렇게 하면 이벤트 시간 워터마크가 운반하는 시간이 감소하지 않도록 보장됩니다. 기본적으로 combineWaitForAllChannels 설정은 false입니다.

  4. [선택] 프레임워크의 WatermarkHandlingStrategy

    WatermarkHandlingStrategy는 사용자 정의 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;
    }
}

더 알아보기 (Learn more)