LabHub
배우기 러닝패스 코스

Apache Flink — 스트림을 엔진으로 돌린다 · 이벤트 시간과 워터마크 · 실습

지연을 바꿔 가며 버려진 행을 센다

LabHub 에서 이어서 보기

목표

이벤트 시간 창이 워터마크에 따라 언제 닫히고 어떤 행을 버리는지, 같은 자료를 배치와 여러 지연의 스트리밍으로 돌려 숫자로 확인한다. CURRENT_WATERMARK 로 늦은 행을 직접 찾아 창이 버리는 행과 비교한다.

왜 중요한가

워터마크가 창을 닫은 뒤 도착한 행은 오류 없이 사라진다. 그래서 "숫자가 조금 모자라다" 는 문제는 로그로 찾을 수 없고, 워터마크가 어떻게 움직이는지 알아야만 설명할 수 있다. 이 실습의 입력은 초 단위 타임스탬프라 워터마크가 레코드마다 즉시 전진하고, 결과가 실행 속도와 무관하게 도착 순서만으로 정해진다. 채점기는 클러스터에 묻지 않는다 — 여러분이 저장한 sql-client 출력을 읽고, 원본 CSV 를 도착 순서대로 흘려 엔진과 같은 규칙(행을 먼저 내보내고 워터마크를 올린다 · window_end ≤ 워터마크 인 창의 행은 버린다)으로 계산한 값과 대조한다.

단계

1. flink-up 으로 클러스터를 띄우고, /root/flink/watermark/ddl.sqlWATERMARK FOR ts AS ts - INTERVAL '5' SECOND 를 둔 events 표와 DESCRIBE events; 를 써서 출력을 /root/flink/watermark/ddl.out 에 저장하세요.
2. 배치 모드로 1분 TUMBLE 창마다 cnt(건수)·total(reading 합)을 내는 /root/flink/watermark/batch.sql 을 돌려 출력을 /root/flink/watermark/batch.out 에 저장하세요.
3. 같은 집계를 스트리밍 모드(지연 5초)로 돌리는 /root/flink/watermark/w5.sql 의 출력을 /root/flink/watermark/w5.out 에 저장하세요.
4. /root/flink/watermark/sweep.sql 에 워터마크가 ts 인 표 events_0ts - INTERVAL '30' SECOND 인 표 events_30 을 만들고, 같은 집계를 events_0events_30 순서로 돌려 출력을 /root/flink/watermark/sweep.out 에 저장하세요.
5. 5초 지연 events 에서 CURRENT_WATERMARK(ts) 가 NULL 이 아니고 ts <= CURRENT_WATERMARK(ts) 인 행의 event_id, ts, wm 을 뽑는 /root/flink/watermark/late.sql 의 출력을 /root/flink/watermark/late.out 에 저장하세요.
6. 늦은 행을 창 앞에서 걸러 낸(CURRENT_WATERMARK(ts) IS NULL OR ts > CURRENT_WATERMARK(ts)) 뒤 같은 1분 집계를 내는 /root/flink/watermark/filtered.sql 의 출력을 /root/flink/watermark/filtered.out 에 저장하세요.
7. 3단계와 같은 집계에 SET 'pipeline.auto-watermark-interval' = '1 h'; 만 더한 /root/flink/watermark/slow.sql 의 출력을 /root/flink/watermark/slow.out 에 저장하세요.
8. /root/flink/watermark/report.jsontotal_rows·dropped_0·dropped_5·dropped_30·late_rows_5·dropped_slow 를 적으세요.

참고

8단계

  1. ts 를 이벤트 시간으로 선언한다
  2. 배치로 기준선을 만든다
  3. 지연 5초 스트리밍 — 버려지는 행
  4. 지연 0초와 30초를 한 번에 비교한다
  5. CURRENT_WATERMARK 로 늦은 행을 뽑는다
  6. 창 앞에서 거르면 더 많이 버린다
  7. 워터마크 간격을 1시간으로 — 전진이 멈춘다
  8. 보고서 — 지연과 버려진 행