워터마크 생성

워터마크 생성 (Generating Watermarks)

이번 섹션에서는 Flink가 이벤트 시간 타임스탬프와 워터마크를 다루기 위해 제공하는 API에 대해 배웁니다. 이벤트 시간, 처리 시간, 수집 시간에 대한 소개는 이벤트 시간 소개를 참고하세요.

출처: 문서

본문

워터마크 전략 소개 (Introduction to Watermark Strategies)

이벤트 시간으로 작업하려면 Flink가 이벤트의 타임스탬프를 알아야 합니다. 즉, 스트림의 각 요소에 이벤트 타임스탬프가 할당되어야 합니다. 이는 보통 TimestampAssigner를 사용해 요소의 특정 필드에서 타임스탬프를 접근/추출하는 방식으로 수행됩니다.

타임스탬프 할당은 이벤트 시간의 진행 상황을 시스템에 알려주는 워터마크 생성과 밀접하게 맞물려 있습니다. 이는 WatermarkGenerator를 지정해 구성할 수 있습니다.

Flink API는 TimestampAssignerWatermarkGenerator를 모두 포함하는 WatermarkStrategy를 기대합니다. 몇 가지 공통 전략이 WatermarkStrategy의 정적 메서드로 즉시 사용 가능하며, 사용자는 필요할 때 자신만의 전략을 만들 수도 있습니다.

참고로 인터페이스는 다음과 같습니다:

public interface WatermarkStrategy<T>
    extends TimestampAssignerSupplier<T>,
            WatermarkGeneratorSupplier<T>{

    /**
     * Instantiates a {@link TimestampAssigner} for assigning timestamps according to this
     * strategy.
     */
    @Override
    TimestampAssigner<T> createTimestampAssigner(TimestampAssignerSupplier.Context context);

    /**
     * Instantiates a WatermarkGenerator that generates watermarks according to this strategy.
     */
    @Override
    WatermarkGenerator<T> createWatermarkGenerator(WatermarkGeneratorSupplier.Context context);
}

언급했듯이 보통 이 인터페이스를 직접 구현하지 않고, 공통 워터마크 전략을 위해 WatermarkStrategy의 정적 헬퍼 메서드를 사용하거나 커스텀 TimestampAssignerWatermarkGenerator를 함께 묶는 데 사용합니다. 예를 들어 유한 지연(bounded-out-of-orderness) 워터마크와 람다 함수를 타임스탬프 할당자로 사용하려면 이렇게 합니다:

Java

WatermarkStrategy
        .<Tuple2<Long, String>>forBoundedOutOfOrderness(Duration.ofSeconds(20))
        .withTimestampAssigner((event, timestamp) -> event.f0);

Python

class FirstElementTimestampAssigner(TimestampAssigner):

    def extract_timestamp(self, value, record_timestamp):
        return value[0]

WatermarkStrategy \
    .for_bounded_out_of_orderness(Duration.of_seconds(20)) \
    .with_timestamp_assigner(FirstElementTimestampAssigner())

TimestampAssigner 지정은 선택 사항이며 대부분의 경우 실제로 지정하고 싶지 않을 것입니다. 예를 들어 Kafka나 Kinesis를 사용할 때 Kafka/Kinesis 레코드에서 타임스탬프를 직접 얻을 수 있습니다.

WatermarkGenerator 인터페이스는 뒤의 Writing WatermarkGenerators 섹션에서 살펴보겠습니다.

주의: 타임스탬프와 워터마크 모두 Java epoch인 1970-01-01T00:00:00Z 이후의 밀리초로 지정됩니다.

워터마크 전략 사용 (Using Watermark Strategies)

Flink 애플리케이션에서 WatermarkStrategy를 사용할 수 있는 두 곳이 있습니다: ① 소스에 직접, ② 비소스 연산 후.

첫 번째 옵션이 선호됩니다. 소스가 워터마킹 로직에서 샤드/파티션/스플릿에 대한 지식을 활용할 수 있기 때문입니다. 소스는 보통 더 세밀한 수준에서 워터마크를 추적할 수 있고 소스가 만들어내는 전체 워터마크가 더 정확해집니다. 소스에 WatermarkStrategy를 직접 지정하려면 보통 소스 특유의 인터페이스를 사용해야 합니다. Kafka 커넥터에서 이것이 어떻게 동작하는지, 그리고 그곳에서 파티션별 워터마킹이 어떻게 동작하는지에 대한 자세한 내용은 Watermark Strategies와 Kafka Connector를 참고하세요.

두 번째 옵션(임의 연산 후 WatermarkStrategy 설정)은 소스에 전략을 직접 설정할 수 없을 때만 사용해야 합니다:

Java

final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

DataStream<MyEvent> stream = env.readFile(
        myFormat, myFilePath, FileProcessingMode.PROCESS_CONTINUOUSLY, 100,
        FilePathFilter.createDefaultFilter(), typeInfo);

DataStream<MyEvent> withTimestampsAndWatermarks = stream
        .filter( event -> event.severity() == WARNING )
        .assignTimestampsAndWatermarks(<watermark strategy>);

withTimestampsAndWatermarks
        .keyBy( (event) -> event.getGroup() )
        .window(TumblingEventTimeWindows.of(Duration.ofSeconds(10)))
        .reduce( (a, b) -> a.add(b) )
        .addSink(...);

Python

env = StreamExecutionEnvironment.get_execution_environment()

# currently read_file is not supported in PyFlink
stream = env \
    .read_text_file(my_file_path, charset) \
    .map(lambda s: MyEvent.from_string(s))

with_timestamp_and_watermarks = stream \
    .filter(lambda e: e.severity() == WARNING) \
    .assign_timestamp_and_watermarks(<watermark strategy>)

with_timestamp_and_watermarks \
    .key_by(lambda e: e.get_group()) \
    .window(TumblingEventTimeWindows.of(Duration.ofSeconds(10))) \
    .reduce(lambda a, b: a.add(b)) \
    .add_sink(...)

이렇게 WatermarkStrategy를 사용하면 스트림을 받아 타임스탬프가 있는 요소와 워터마크가 있는 새 스트림을 생성합니다. 원래 스트림에 이미 타임스탬프나 워터마크가 있었다면 타임스탬프 할당자가 그것을 덮어씁니다.

유휴 소스 처리 (Dealing With Idle Sources)

입력 스플릿/파티션/샤드 중 하나가 한동안 이벤트를 전달하지 않으면 WatermarkGenerator도 워터마크의 기반이 될 새 정보를 얻지 못합니다. 이를 유휴 입력(idle input) 또는 유휴 소스(idle source)라고 합니다. 일부 파티션은 여전히 이벤트를 전달할 수 있기 때문에 이것은 문제입니다. 그 경우 워터마크가 모든 서로 다른 병렬 워터마크의 최솟값으로 계산되므로, 워터마크가 지연(hold back)됩니다.

이를 처리하기 위해 유휴 상태를 감지하고 입력을 유휴로 표시하는 WatermarkStrategy를 사용할 수 있습니다. WatermarkStrategy는 이를 위한 편리한 헬퍼를 제공합니다:

Java

WatermarkStrategy
        .<Tuple2<Long, String>>forBoundedOutOfOrderness(Duration.ofSeconds(20))
        .withIdleness(Duration.ofMinutes(1));

Python

WatermarkStrategy \
    .for_bounded_out_of_orderness(Duration.of_seconds(20)) \
    .with_idleness(Duration.of_minutes(1))

워터마크 정렬 (Watermark alignment)

앞 문단에서는 스플릿/파티션/샤드 또는 소스가 유휴 상태가 되어 워터마크 증가를 막을 수 있는 상황을 논의했습니다. 반대쪽 스펙트럼으로, 스플릿/파티션/샤드 또는 소스가 레코드를 매우 빠르게 처리해 다른 것보다 상대적으로 빠르게 워터마크를 올릴 수 있습니다. 이것 자체로는 문제가 아닙니다. 그러나 워터마크를 사용해 데이터를 내보내는 다운스트림 연산자에게는 실제로 문제가 될 수 있습니다.

이 경우 유휴 소스와는 반대로, 그러한 다운스트림 연산자(예: 집계에 대한 윈도우 조인)의 워터마크는 진행될 수 있습니다. 그러나 그러한 연산자는 모든 입력으로부터의 최소 워터마크가 뒤처지는 입력에 의해 지연되기 때문에, 빠른 입력에서 오는 과도한 양의 데이터를 버퍼링해야 할 필요가 있습니다. 빠른 입력이 내보낸 모든 레코드는 그러한 다운스트림 연산자 상태에 버퍼링되어야 하며, 이는 연산자 상태의 통제 불가능한 성장으로 이어질 수 있습니다.

이 문제를 해결하기 위해 워터마크 정렬(watermark alignment)을 활성화할 수 있습니다. 이것은 어떤 소스/스플릿/샤드/파티션도 나머지보다 너무 앞서 워터마크를 올리지 못하게 합니다. 각 소스에 대해 개별적으로 정렬을 활성화할 수 있습니다:

Java

WatermarkStrategy
        .<Tuple2<Long, String>>forBoundedOutOfOrderness(Duration.ofSeconds(20))
        .withWatermarkAlignment("alignment-group-1", Duration.ofSeconds(20), Duration.ofSeconds(1));

Python

WatermarkStrategy \
    .for_bounded_out_of_orderness(Duration.of_seconds(20)) \
    .with_watermark_alignment("alignment-group-1", Duration.of_seconds(20), Duration.of_seconds(1))

참고: 워터마크 정렬은 FLIP-27 소스에서만 활성화할 수 있습니다. 레거시 소스나 소스 이후 DataStream#assignTimestampsAndWatermarks로 적용한 경우에는 동작하지 않습니다.

정렬을 활성화할 때 소스가 속할 그룹을 Flink에 알려야 합니다. 이것은 그 그룹을 공유하는 모든 소스를 묶는 레이블(예: alignment-group-1)을 제공해 수행합니다. 또한 해당 그룹에 속한 모든 소스의 현재 최소 워터마크에서의 최대 드리프트(drift)를 알려야 합니다. 세 번째 파라미터는 현재 최대 워터마크를 얼마나 자주 갱신할지 설명합니다. 잦은 갱신의 단점은 TM과 JM 사이에 더 많은 RPC 메시지가 오간다는 것입니다.

정렬을 달성하기 위해 Flink는 너무 먼 미래의 워터마크를 생성한 소스/태스크에서 소비를 일시 중지합니다. 그동안 결합된 워터마크를 앞으로 움직일 수 있는 다른 소스/태스크에서 레코드를 계속 읽어 더 빠른 쪽의 차단을 해제합니다.

참고: Flink 1.17부터 FLIP-27 소스 프레임워크에서 스플릿 레벨 워터마크 정렬을 지원합니다. 소스 커넥터는 같은 태스크에서 스플릿/파티션/샤드를 정렬할 수 있도록 스플릿을 재개·일시 중지하는 인터페이스를 구현해야 합니다. 일시 중지·재개 인터페이스의 자세한 내용은 Source API에서 확인할 수 있습니다.

1.15.x에서 1.16.x 사이 버전(포함)의 Flink에서 업그레이드하는 경우 pipeline.watermark-alignment.allow-unaligned-source-splits를 true로 설정해 스플릿 레벨 정렬을 비활성화할 수 있습니다. 또한 소스가 런타임에 UnsupportedOperationException을 던지는지 또는 javadocs를 읽어 스플릿 레벨 정렬을 지원하는지 알 수 있습니다. 그런 경우 치명적 예외를 피하기 위해 스플릿 레벨 워터마크 정렬을 비활성화하는 것이 바람직합니다.

플래그를 true로 설정하면 워터마크 정렬은 스플릿/샤드/파티션 수가 소스 연산자의 병렬도와 같을 때만 제대로 동작합니다. 이렇게 하면 각 서브태스크에 단일 작업 단위가 할당됩니다. 반면 서로 다른 속도로 워터마크를 생성하고 같은 태스크에 할당된 두 개의 Kafka 파티션이 있다면 워터마크가 예상대로 동작하지 않을 수 있습니다. 다행히 최악의 경우에도 기본 정렬은 정렬이 전혀 없는 것보다 나쁘게 동작하지 않아야 합니다.

또한 Flink는 같은 소스 및/또는 서로 다른 소스의 태스크 간 정렬도 지원합니다. 이는 서로 다른 속도로 워터마크를 생성하는 두 개의 서로 다른 소스(예: Kafka와 File)가 있을 때 유용합니다.

WatermarkGenerator 작성 (Writing WatermarkGenerators)

TimestampAssigner는 이벤트에서 필드를 추출하는 단순한 함수이므로 자세히 볼 필요가 없습니다. 반면 WatermarkGenerator는 작성이 조금 더 복잡하며, 다음 두 섹션에서 그 방법을 살펴보겠습니다. 이것이 WatermarkGenerator 인터페이스입니다:

/**
 * The {@code WatermarkGenerator} generates watermarks either based on events or
 * periodically (in a fixed interval).
 *
 * <p><b>Note:</b> This WatermarkGenerator subsumes the previous distinction between the
 * {@code AssignerWithPunctuatedWatermarks} and the {@code AssignerWithPeriodicWatermarks}.
 */
@Public
public interface WatermarkGenerator<T> {

    /**
     * Called for every event, allows the watermark generator to examine
     * and remember the event timestamps, or to emit a watermark based on
     * the event itself.
     */
    void onEvent(T event, long eventTimestamp, WatermarkOutput output);

    /**
     * Called periodically, and might emit a new watermark, or not.
     *
     * <p>The interval in which this method is called and Watermarks
     * are generated depends on {@link ExecutionConfig#getAutoWatermarkInterval()}.
     */
    void onPeriodicEmit(WatermarkOutput output);
}

워터마크 생성에는 주기적(periodic)과 구두점(punctuated) 두 가지 스타일이 있습니다.

주기적 생성기는 보통 onEvent()로 들어오는 이벤트를 관찰하고, 프레임워크가 onPeriodicEmit()을 호출할 때 워터마크를 내보냅니다.

구두점 생성기는 onEvent()에서 이벤트를 보고 스트림에서 워터마크 정보를 담은 특수 마커 이벤트나 구두점(punctuation)을 기다립니다. 그러한 이벤트를 보면 즉시 워터마크를 내보냅니다. 보통 구두점 생성기는 onPeriodicEmit()에서 워터마크를 내보내지 않습니다.

다음에서 각 스타일의 생성기를 구현하는 방법을 살펴보겠습니다.

주기적 WatermarkGenerator 작성 (Writing a Periodic WatermarkGenerator)

주기적 생성기는 스트림 이벤트를 관찰하고 주기적으로(스트림 요소에 따라 또는 순전히 처리 시간에 기반해) 워터마크를 생성합니다.

워터마크가 생성될 간격(매 n 밀리초)은 ExecutionConfig.setAutoWatermarkInterval(...)로 정의됩니다. 생성기의 onPeriodicEmit() 메서드가 매번 호출되며, 반환된 워터마크가 null이 아니고 이전 워터마크보다 크면 새 워터마크가 내보집니다.

여기 주기적 워터마크 생성을 사용하는 워터마크 생성기의 두 가지 간단한 예를 보여줍니다. Flink에는 아래에 보이는 BoundedOutOfOrdernessGenerator와 유사하게 동작하는 WatermarkGeneratorBoundedOutOfOrdernessWatermarks가 함께 제공됩니다. 그 사용법은 여기에서 읽을 수 있습니다.

Java

/**
 * This generator generates watermarks assuming that elements arrive out of order,
 * but only to a certain degree. The latest elements for a certain timestamp t will arrive
 * at most n milliseconds after the earliest elements for timestamp t.
 */
public class BoundedOutOfOrdernessGenerator implements WatermarkGenerator<MyEvent> {

    private final long maxOutOfOrderness = 3500; // 3.5 seconds

    private long currentMaxTimestamp;

    @Override
    public void onEvent(MyEvent event, long eventTimestamp, WatermarkOutput output) {
        currentMaxTimestamp = Math.max(currentMaxTimestamp, eventTimestamp);
    }

    @Override
    public void onPeriodicEmit(WatermarkOutput output) {
        // emit the watermark as current highest timestamp minus the out-of-orderness bound
        output.emitWatermark(new Watermark(currentMaxTimestamp - maxOutOfOrderness - 1));
    }

}

/**
 * This generator generates watermarks that are lagging behind processing time
 * by a fixed amount. It assumes that elements arrive in Flink after a bounded delay.
 */
public class TimeLagWatermarkGenerator implements WatermarkGenerator<MyEvent> {

    private final long maxTimeLag = 5000; // 5 seconds

    @Override
    public void onEvent(MyEvent event, long eventTimestamp, WatermarkOutput output) {
        // don't need to do anything because we work on processing time
    }

    @Override
    public void onPeriodicEmit(WatermarkOutput output) {
        output.emitWatermark(new Watermark(System.currentTimeMillis() - maxTimeLag));
    }
}

Python

Still not supported in Python API.

구두점 WatermarkGenerator 작성 (Writing a Punctuated WatermarkGenerator)

구두점 워터마크 생성기는 이벤트 스트림을 관찰하고 워터마크 정보를 담은 특수 요소를 볼 때마다 워터마크를 내보냅니다.

이벤트가 특정 마커를 담고 있음을 나타낼 때 워터마크를 내보내는 구두점 생성기를 구현하는 방법은 다음과 같습니다:

Java

public class PunctuatedAssigner implements WatermarkGenerator<MyEvent> {

    @Override
    public void onEvent(MyEvent event, long eventTimestamp, WatermarkOutput output) {
        if (event.hasWatermarkMarker()) {
            output.emitWatermark(new Watermark(event.getWatermarkTimestamp()));
        }
    }

    @Override
    public void onPeriodicEmit(WatermarkOutput output) {
        // don't need to do anything because we emit in reaction to events above
    }
}

Python

Still not supported in Python API.

참고: 모든 단일 이벤트에서 워터마크를 생성하는 것도 가능합니다. 그러나 각 워터마크가 다운스트림에서 일부 계산을 일으키므로 과도한 워터마크 수는 성능을 저하시킵니다.

워터마크 전략과 Kafka 커넥터 (Watermark Strategies and the Kafka Connector)

Apache Kafka를 데이터 소스로 사용할 때 각 Kafka 파티션은 단순한 이벤트 시간 패턴(오름차순 타임스탬프 또는 유한 지연)을 가질 수 있습니다. 그러나 Kafka에서 스트림을 소비할 때 여러 파티션이 종종 병렬로 소비되어, 파티션들의 이벤트를 인터리브하고 파티션별 패턴을 파괴합니다(이것은 Kafka 컨슈머 클라이언트가 동작하는 방식에 내재된 것입니다).

그 경우 Flink의 Kafka-파티션 인식 워터마크 생성을 사용할 수 있습니다. 이 기능을 사용하면 워터마크가 Kafka 컨슈머 내부에서 Kafka 파티션별로 생성되며, 파티션별 워터마크는 스트림 셔플에서 워터마크가 병합되는 것과 같은 방식으로 병합됩니다.

예를 들어 이벤트 타임스탬프가 Kafka 파티션별로 엄격히 오름차순이면, 오름차순 타임스탬프 워터마크 생성기로 파티션별 워터마크를 생성하면 완벽한 전체 워터마크가 됩니다. 예시에서 TimestampAssigner를 제공하지 않으며, 대신 Kafka 레코드 자체의 타임스탬프가 사용됨에 유의하세요.

아래 그림은 Kafka 파티션별 워터마크 생성의 사용 방법과 그 경우 워터마크가 스트리밍 데이터플로우를 통해 어떻게 전파되는지 보여줍니다.

Java

KafkaSource<String> kafkaSource = KafkaSource.<String>builder()
    .setBootstrapServers(brokers)
    .setTopics("my-topic")
    .setGroupId("my-group")
    .setStartingOffsets(OffsetsInitializer.earliest())
    .setValueOnlyDeserializer(new SimpleStringSchema())
    .build();

DataStream<String> stream = env.fromSource(
    kafkaSource, WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(20)), "mySource");

Python

kafka_source = KafkaSource.builder()
    .set_bootstrap_servers(brokers)
    .set_topics("my-topic")
    .set_group_id("my-group")
    .set_starting_offsets(KafkaOffsetsInitializer.earliest())
    .set_value_only_deserializer(SimpleStringSchema())
    .build()

stream = env.from_source(
    source=kafka_source,
    watermark_strategy=WatermarkStrategy.for_bounded_out_of_orderness(Duration.of_seconds(20)),
    source_name="kafka_source")

연산자가 워터마크를 처리하는 방법 (How Operators Process Watermarks)

일반적인 규칙으로, 연산자는 주어진 워터마크를 하류로 전달하기 전에 완전히 처리해야 합니다. 예를 들어 WindowOperator는 먼저 트리거되어야 할 모든 윈도우를 평가하고, 워터마크가 유발한 모든 출력을 생성한 후에야 워터마크 자체를 하류로 보냅니다. 즉, 워터마크 발생으로 인해 생성된 모든 요소는 워터마크보다 먼저 내보내집니다.

같은 규칙이 TwoInputStreamOperator에도 적용됩니다. 그러나 이 경우 연산자의 현재 워터마크는 두 입력의 최솟값으로 정의됩니다.

이 동작의 세부 사항은 OneInputStreamOperator#processWatermark, TwoInputStreamOperator#processWatermark1, TwoInputStreamOperator#processWatermark2 메서드의 구현으로 정의됩니다.

더 이상 사용되지 않는 AssignerWithPeriodicWatermarks와 AssignerWithPunctuatedWatermarks (The Deprecated AssignerWithPeriodicWatermarks and AssignerWithPunctuatedWatermarks)

현재의 WatermarkStrategy, TimestampAssigner, WatermarkGenerator 추상화가 도입되기 전에 Flink는 AssignerWithPeriodicWatermarksAssignerWithPunctuatedWatermarks를 사용했습니다. 여전히 API에서 볼 수 있지만, 새 인터페이스가 관심사의 더 명확한 분리를 제공하고 주기적·구두점 워터마크 생성 스타일을 통합하므로 새 인터페이스를 사용하는 것을 권장합니다.

더 알아보기 (Learn more)