상위: Flink
요약
Flink의 Watermark 메커니즘을 다룹니다. 이벤트 시간과 처리 시간의 차이를 해결하고, 늦게 도착하는 데이터를 처리하는 방법을 학습합니다. Watermark 생성 전략과 Allowed Lateness를 통한 지연 데이터 처리 방법을 알아봅니다.
Watermarks
- 타임스탬프 메타데이터로서, 현재까지 도착한 이벤트 중 가장 큰 이벤트 시간에 대한 정보.
- 예: "현재 시각(t)까지의 이벤트는 모두 도착했다고 간주하겠다"는 마커를 주기적으로 스트림과 함께 전송.
- 윈도우 연산자는 이 워터마크를 참고하여 윈도우를 닫을 시점을 결정.

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

필요성
- 실시간 스트림 처리에서 이벤트의 실제 발생 시각(Event Time)과 시스템 처리 시각(Processing Time)이 다를 수 있음.
- 네트워크 지연, 시스템 부하 등으로 인해 이벤트가 순서대로 도착하지 않을 수 있음.
- 아무 대책 없이 이벤트 시간 기준 윈도우를 사용하면, 이미 닫힌 윈도우에 늦게 도착한 이벤트는 반영되지 못해 데이터 손실 발생 가능.
Event Time vs. Processing Time
- Event Time
- 데이터가 실제 현실에서 발생한 시각.
- Event Time 기반 윈도우 처리 시 발생 시각을 정확히 반영하여 논리적으로 일관된 결과 확보 가능.
- 워터마크 설정 필수 (Flink가 시간 진행을 알 수 있도록)
- Processing Time
- Flink 태스크가 이벤트를 처리하는 현재 시스템 시간.
- 구현이 간단하고 실시간성은 높지만, 외부 요인으로 순서가 뒤바뀌거나 지연된 이벤트는 고려할 수 없음.
- 정확도는 낮아질 수 있음.
작동 원리
- 워터마크는 일반적으로
**현재 이벤트 시간 - 허용 지연 시간**으로 생성 - 해당 워터마크 시간보다 이전에 끝나는 윈도우들은 더 이상 이벤트가 없다고 판단하고 종료


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

- Allowed Lateness는 Watermark delay 이후 적용

예제 코드 (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))
관련 개념
- Window Implementation - 윈도우 구현 방법
- Triggers - 윈도우 트리거 메커니즘
- Evictors - 윈도우 데이터 제거