내장 워터마크 생성기

내장 워터마크 생성기 (Builtin Watermark Generators)

Flink가 기본 제공하는 타임스탬프 할당자를 소개해요. 단조 증가 타임스탬프와 고정 지연량 워터마크 생성이라는 두 내장 워터마크 생성 방식을 다룹니다.

출처: Builtin Watermark Generators

본문

Generating Watermarks에서 설명했듯이 Flink는 프로그래머가 직접 타임스탬프를 할당하고 워터마크를 생성할 수 있는 추상화를 제공해요. 구체적으로는 WatermarkGenerator 인터페이스를 구현하면 돼요.

이런 작업의 프로그래밍 부담을 더 줄이기 위해 Flink는 미리 구현된 타임스탬프 할당자를 함께 제공해요. 이 섹션은 그것들의 목록을 제공해요. 기본 제공 기능 외에도 그 구현은 커스텀 구현의 예시로 사용될 수 있어요.

Monotonously Increasing Timestamps

주기적(periodic) 워터마크 생성의 가장 단순한 특수 경우는 주어진 소스 태스크가 보는 타임스탬프가 오름차순으로 발생할 때예요. 이 경우 이전 타임스탬프는 도착하지 않으므로 현재 타임스탬프가 항상 워터마크 역할을 할 수 있어요.

타임스탬프가 병렬 데이터 소스 태스크별로 오름차순이면 충분하다는 점에 유의하세요. 예를 들어 특정 설정에서 하나의 Kafka 파티션을 하나의 병렬 데이터 소스 인스턴스가 읽는다면 각 Kafka 파티션 안에서만 타임스탬프가 오름차순이면 충분해요. Flink의 워터마크 병합 메커니즘은 병렬 스트림이 셔플·유니온·연결·병합될 때마다 올바른 워터마크를 생성해요.

WatermarkStrategy.forMonotonousTimestamps();
WatermarkStrategy.for_monotonous_timestamps()

Fixed Amount of Lateness

주기적 워터마크 생성의 또 다른 예는 워터마크가 스트림에서 본 최대 (이벤트 타임) 타임스탬프보다 고정된 시간만큼 뒤처지는 경우예요. 이 경우는 스트림에서 발생할 수 있는 최대 지연(lateness)을 미리 알고 있는 시나리오를 다뤄요. 예를 들어 테스트를 위해 고정 기간 내에 퍼진 타임스탬프를 가진 요소를 포함하는 커스텀 소스를 만들 때가 그렇죠. 이런 경우를 위해 Flink는 BoundedOutOfOrdernessWatermarks 생성기를 제공하는데, 인자로 maxOutOfOrderness, 즉 주어진 윈도우의 최종 결과를 계산할 때 무시되기 전에 요소가 늦어질 수 있는 최대 시간을 받아요. 지연(lateness)은 t_w - t의 결과에 해당하며, 여기서 t는 요소의 (이벤트 타임) 타임스탬프이고 t_w는 이전 워터마크의 타임스탬프예요. lateness > 0이면 요소는 늦은 것으로 간주되고 기본적으로 해당 윈도우의 잡 결과를 계산할 때 무시돼요. 늦은 요소를 다루는 자세한 내용은 allowed lateness 문서를 참고하세요.

WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(10));
WatermarkStrategy.for_bounded_out_of_orderness(Duration.of_seconds(10))

더 알아보기 (Learn more)