Apache Flink — 스트림을 엔진으로 돌린다 · 이벤트 시간과 워터마크 · 실습
지연을 바꿔 가며 버려진 행을 센다
목표
이벤트 시간 창이 워터마크에 따라 언제 닫히고 어떤 행을 버리는지, 같은 자료를 배치와 여러 지연의 스트리밍으로 돌려 숫자로 확인한다. CURRENT_WATERMARK 로 늦은 행을 직접 찾아 창이 버리는 행과 비교한다.
왜 중요한가
워터마크가 창을 닫은 뒤 도착한 행은 오류 없이 사라진다. 그래서 "숫자가 조금 모자라다" 는 문제는 로그로 찾을 수 없고, 워터마크가 어떻게 움직이는지 알아야만 설명할 수 있다. 이 실습의 입력은 초 단위 타임스탬프라 워터마크가 레코드마다 즉시 전진하고, 결과가 실행 속도와 무관하게 도착 순서만으로 정해진다. 채점기는 클러스터에 묻지 않는다 — 여러분이 저장한 sql-client 출력을 읽고, 원본 CSV 를 도착 순서대로 흘려 엔진과 같은 규칙(행을 먼저 내보내고 워터마크를 올린다 · window_end ≤ 워터마크 인 창의 행은 버린다)으로 계산한 값과 대조한다.
단계
1. flink-up 으로 클러스터를 띄우고, /root/flink/watermark/ddl.sql 에 WATERMARK 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_0 과 ts - INTERVAL '30' SECOND 인 표 events_30 을 만들고, 같은 집계를 events_0 → events_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.json 에 total_rows·dropped_0·dropped_5·dropped_30·late_rows_5·dropped_slow 를 적으세요.
참고
- 원본 열:
event_id BIGINT, sensor STRING, reading INT, ts TIMESTAMP(3)(머리글 없는 CSV, 파일/opt/lab/fixtures/data/watermark_events.csv, 파일 순서 = 도착 순서). - 창 집계 모양:
SELECT window_start, window_end, COUNT(*) AS cnt, SUM(reading) AS total FROM TUMBLE(TABLE 표, DESCRIPTOR(ts), INTERVAL '1' MINUTE) GROUP BY window_start, window_end; - 한 SQL 파일에 SELECT 가 여럿이면 차례로 잡이 돌고, 출력에 결과 표가 순서대로 찍힙니다.
- 흔한 실수: 배치 모드는 워터마크를 쓰지 않습니다. 늦은 행을 보려면 스트리밍으로 돌려야 합니다. 병렬도는 기본값 1 그대로 두세요 — 입력이 여럿이면 워터마크는 그중 최솟값을 따릅니다.
- 흔한 실수: 첫 행에서는 워터마크가 아직 없어
CURRENT_WATERMARK(ts)가 NULL 입니다. NULL 과의 비교는 참이 아니라 WHERE 를 통과하지 못하므로, 6단계의 거르기 식에서IS NULL조건을 빼면 첫 행까지 버려집니다. - 공식 문서: [Timely Stream Processing](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/concepts/time/) · [CREATE — WATERMARK](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/sql/reference/ddl/create/) · [Time Attributes](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/concepts/sql-table-concepts/time_attributes/) · [Windowing TVF](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/sql/reference/queries/window-tvf/) · [Built-in Functions](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/sql/functions/built-in-functions/) · [Configuration](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/deployment/config/)
8단계
- ts 를 이벤트 시간으로 선언한다
- 배치로 기준선을 만든다
- 지연 5초 스트리밍 — 버려지는 행
- 지연 0초와 30초를 한 번에 비교한다
- CURRENT_WATERMARK 로 늦은 행을 뽑는다
- 창 앞에서 거르면 더 많이 버린다
- 워터마크 간격을 1시간으로 — 전진이 멈춘다
- 보고서 — 지연과 버려진 행