전체 그래프
Flink

Watermarks

data-engineeringflinkwatermarksevent-timeprocessing-timelate-elements

상위: Flink

요약

Flink의 Watermark 메커니즘을 다룹니다. 이벤트 시간과 처리 시간의 차이를 해결하고, 늦게 도착하는 데이터를 처리하는 방법을 학습합니다. Watermark 생성 전략과 Allowed Lateness를 통한 지연 데이터 처리 방법을 알아봅니다.

Watermarks

  • 타임스탬프 메타데이터로서, 현재까지 도착한 이벤트 중 가장 큰 이벤트 시간에 대한 정보.
  • 예: "현재 시각(t)까지의 이벤트는 모두 도착했다고 간주하겠다"는 마커를 주기적으로 스트림과 함께 전송.
  • 윈도우 연산자는 이 워터마크를 참고하여 윈도우를 닫을 시점을 결정.

flink-33.png

  • Window Operator에 Watermark가 유입되고 나면 Window Operator는 더 이상 1h 13m보다 과거의 Event는 유입되지 않는다고 판단
  • Watermark가 명시한 시간보다 과거의 Event들은 모두 Purge 됨

flink-34.png

필요성

  • 실시간 스트림 처리에서 이벤트의 실제 발생 시각(Event Time)과 시스템 처리 시각(Processing Time)이 다를 수 있음.
  • 네트워크 지연, 시스템 부하 등으로 인해 이벤트가 순서대로 도착하지 않을 수 있음.
  • 아무 대책 없이 이벤트 시간 기준 윈도우를 사용하면, 이미 닫힌 윈도우에 늦게 도착한 이벤트는 반영되지 못해 데이터 손실 발생 가능.

Event Time vs. Processing Time

  • Event Time
    • 데이터가 실제 현실에서 발생한 시각.
    • Event Time 기반 윈도우 처리 시 발생 시각을 정확히 반영하여 논리적으로 일관된 결과 확보 가능.
    • 워터마크 설정 필수 (Flink가 시간 진행을 알 수 있도록)
  • Processing Time
    • Flink 태스크가 이벤트를 처리하는 현재 시스템 시간.
    • 구현이 간단하고 실시간성은 높지만, 외부 요인으로 순서가 뒤바뀌거나 지연된 이벤트는 고려할 수 없음.
    • 정확도는 낮아질 수 있음.

작동 원리

  • 워터마크는 일반적으로 **현재 이벤트 시간 - 허용 지연 시간**으로 생성
  • 해당 워터마크 시간보다 이전에 끝나는 윈도우들은 더 이상 이벤트가 없다고 판단하고 종료

flink-35.png

flink-36.png

Late Elements와 Allowed Lateness

  • **Allowed Lateness(Duration)**을 설정하면 윈도우가 닫힌 후에도 일정 시간 동안 늦은 이벤트를 수용 가능.
  • 설정된 기간 내 도착한 이벤트는 윈도우를 다시 열고 계산에 반영함.

flink-37.png

  • Allowed Lateness는 Watermark delay 이후 적용

flink-38.png

예제 코드 (Watermark + Allowed Lateness)

  • 1초까지 out-of-order를 허용하는 워터마크
  • 2초의 이벤트 시간 텀블링 윈도우에 대해 최대 2초의 허용 지연 설정

TimestampAssigner 정의

class CustomTimestampAssigner(TimestampAssigner):
    def extract_timestamp(self, element, record_timestamp):
        return element[1]

ProcessFunction 정의

class PrintWatermarkProcessFunction(ProcessFunction):
    def process_element(self, value, ctx):
        watermark = ctx.timer_service().current_watermark()
        print(f"Event: {value}, Current Watermark: {watermark}")
        yield value

환경 및 데이터 생성

env = StreamExecutionEnvironment.get_execution_environment()
env.set_parallelism(1)

data = [
    (1, int(now.timestamp())),
    (2, int((now + timedelta(milliseconds=1000)).timestamp())),
    (3, int((now + timedelta(milliseconds=2000)).timestamp())),
    ...
]
source = env.from_collection(data, type_info=Types.TUPLE([Types.INT(), Types.LONG()]))

워터마크 전략 및 윈도우 처리 설정

# 1초 out-of-order 허용 워터마크 설정
watermark_strategy = WatermarkStrategy.for_bounded_out_of_orderness(
		Duration.of_seconds(1)) \
    .with_timestamp_assigner(CustomTimestampAssigner()
)

watermarked_stream = source.assign_timestamps_and_watermarks(watermark_strategy)

# 2초 allowed lateness 설정한 윈도우 처리
windowed_stream = watermarked_stream.key_by(
		lambda x: x[0]) \
    .window(TumblingEventTimeWindows.of(Time.seconds(2))) \
    .allowed_lateness(Duration.of_seconds(2)
)

# 처리  출력
processed_stream = watermarked_stream.process(PrintWatermarkProcessFunction())
processed_stream.print()

env.execute("Watermark with Allowed Lateness Example")

출력 :

Event: (1, 1683512335), Current Watermark: 1683512335  
Event: (2, 1683512336), Current Watermark: 1683512336  
Event: (3, 1683512337), Current Watermark: 1683512337  
Event: (4, 1683512338), Current Watermark: 1683512338  
Event: (5, 1683512339), Current Watermark: 1683512339  

Watermark 생성 전략

1. 고정 지연 Watermark

# 5초 고정 지연
watermark_strategy = WatermarkStrategy.for_bounded_out_of_orderness(
    Duration.of_seconds(5)
)

2. 적응형 Watermark

class AdaptiveWatermarkStrategy(WatermarkStrategy):
    def __init__(self, max_lateness):
        self.max_lateness = max_lateness
        self.observed_lateness = 0

    def create_watermark_generator(self, context):
        return AdaptiveWatermarkGenerator(self.max_lateness, self.observed_lateness)

class AdaptiveWatermarkGenerator(WatermarkGenerator):
    def __init__(self, max_lateness, observed_lateness):
        self.max_lateness = max_lateness
        self.observed_lateness = observed_lateness

    def on_event(self, event, event_timestamp, output):
        # 지연 관찰  조정
        current_lateness = time.time() - event_timestamp
        if current_lateness > self.observed_lateness:
            self.observed_lateness = current_lateness
        
        # 적응형 지연 시간 계산
        adaptive_delay = max(self.max_lateness, self.observed_lateness * 1.2)
        watermark = event_timestamp - adaptive_delay
        output.emit_watermark(Watermark(watermark))

3. 주기적 Watermark

class PeriodicWatermarkStrategy(WatermarkStrategy):
    def create_watermark_generator(self, context):
        return PeriodicWatermarkGenerator(interval=1000)  # 1초마다

class PeriodicWatermarkGenerator(WatermarkGenerator):
    def __init__(self, interval):
        self.interval = interval
        self.last_watermark = 0

    def on_periodic_emit(self, output):
        current_time = time.time() * 1000
        if current_time - self.last_watermark >= self.interval:
            watermark = current_time - 5000  # 5초 지연
            output.emit_watermark(Watermark(watermark))
            self.last_watermark = current_time

지연 데이터 처리

1. Side Output 활용

class LateDataHandler(ProcessFunction):
    def __init__(self):
        self.late_data_output = OutputTag("late-data")

    def process_element(self, value, ctx):
        current_watermark = ctx.timer_service().current_watermark()
        event_time = value[1]
        
        if event_time < current_watermark:
            # 지연된 데이터를 Side Output으로 전송
            ctx.output(self.late_data_output, value)
        else:
            # 정상 처리
            yield value

# Side Output 처리
late_data_stream = processed_stream.get_side_output(
    OutputTag("late-data")
)
late_data_stream.print()

2. 지연 데이터 통계

class LateDataStatistics(ProcessFunction):
    def __init__(self):
        self.late_count = 0
        self.total_count = 0

    def process_element(self, value, ctx):
        self.total_count += 1
        current_watermark = ctx.timer_service().current_watermark()
        event_time = value[1]
        
        if event_time < current_watermark:
            self.late_count += 1
            print(f"Late data detected: {self.late_count}/{self.total_count} "
                  f"({self.late_count/self.total_count*100:.2f}%)")
        
        yield value

Watermark 최적화

1. 적절한 지연 시간 설정

# 너무 작은 지연: 데이터 손실 위험
watermark_strategy = WatermarkStrategy.for_bounded_out_of_orderness(
    Duration.of_milliseconds(100)  # 너무 작음
)

# 너무  지연: 메모리 사용량 증가
watermark_strategy = WatermarkStrategy.for_bounded_out_of_orderness(
    Duration.of_hours(1)  # 너무 
)

# 적절한 지연 시간
watermark_strategy = WatermarkStrategy.for_bounded_out_of_orderness(
    Duration.of_seconds(5)  # 적절함
)

2. 메모리 관리

class MemoryAwareWatermarkStrategy(WatermarkStrategy):
    def create_watermark_generator(self, context):
        return MemoryAwareWatermarkGenerator(
            max_memory_usage=1024 * 1024,  # 1MB
            cleanup_interval=60000  # 1분
        )

class MemoryAwareWatermarkGenerator(WatermarkGenerator):
    def __init__(self, max_memory_usage, cleanup_interval):
        self.max_memory_usage = max_memory_usage
        self.cleanup_interval = cleanup_interval
        self.buffered_events = []
        self.last_cleanup = time.time()

    def on_event(self, event, event_timestamp, output):
        # 메모리 사용량 체크
        if self.get_memory_usage() > self.max_memory_usage:
            self.cleanup_old_events()
        
        # 정기적 정리
        if time.time() - self.last_cleanup > self.cleanup_interval:
            self.cleanup_old_events()
            self.last_cleanup = time.time()
        
        # 워터마크 생성
        watermark = event_timestamp - 5000
        output.emit_watermark(Watermark(watermark))

관련 개념