Apache Flink — ストリームを本物のエンジンで動かす
遅延を変えながら捨てられた行を数える
한국어 원문으로 표시합니다.
목표
이벤트 시간 창이 워터마크에 따라 언제 닫히고 어떤 행을 버리는지, 같은 자료를 배치와 여러 지연의 스트리밍으로 돌려 숫자로 확인한다. CURRENT_WATERMARK 로 늦은 행을 직접 찾아 창이 버리는 행과 비교한다.
왜 중요한가
워터마크가 창을 닫은 뒤 도착한 행은 오류 없이 사라진다. 그래서 "숫자가 조금 모자라다" 는 문제는 로그로 찾을 수 없고, 워터마크가 어떻게 움직이는지 알아야만 설명할 수 있다. 이 실습의 입력은 초 단위 타임스탬프라 워터마크가 레코드마다 즉시 전진하고, 결과가 실행 속도와 무관하게 도착 순서만으로 정해진다. 채점기는 클러스터에 묻지 않는다 — 여러분이 저장한 sql-client 출력을 읽고, 원본 CSV 를 도착 순서대로 흘려 엔진과 같은 규칙(행을 먼저 내보내고 워터마크를 올린다 · window_end ≤ 워터마크 인 창의 행은 버린다)으로 계산한 값과 대조한다.
단계
flink-up으로 클러스터를 띄우고, /root/flink/watermark/ddl.sql 에WATERMARK FOR ts AS ts - INTERVAL '5' SECOND를 둔events표와DESCRIBE events;를 써서 출력을 /root/flink/watermark/ddl.out 에 저장하세요.- 배치 모드로 1분
TUMBLE창마다cnt(건수)·total(reading 합)을 내는 /root/flink/watermark/batch.sql 을 돌려 출력을 /root/flink/watermark/batch.out 에 저장하세요. - 같은 집계를 스트리밍 모드(지연 5초)로 돌리는 /root/flink/watermark/w5.sql 의 출력을 /root/flink/watermark/w5.out 에 저장하세요.
- /root/flink/watermark/sweep.sql 에 워터마크가
ts인 표events_0과ts - INTERVAL '30' SECOND인 표events_30을 만들고, 같은 집계를events_0→events_30순서로 돌려 출력을 /root/flink/watermark/sweep.out 에 저장하세요. - 5초 지연
events에서CURRENT_WATERMARK(ts)가 NULL 이 아니고ts <= CURRENT_WATERMARK(ts)인 행의event_id, ts, wm을 뽑는 /root/flink/watermark/late.sql 의 출력을 /root/flink/watermark/late.out 에 저장하세요. - 늦은 행을 창 앞에서 걸러 낸(
CURRENT_WATERMARK(ts) IS NULL OR ts > CURRENT_WATERMARK(ts)) 뒤 같은 1분 집계를 내는 /root/flink/watermark/filtered.sql 의 출력을 /root/flink/watermark/filtered.out 에 저장하세요. - 3단계와 같은 집계에
SET 'pipeline.auto-watermark-interval' = '1 h';만 더한 /root/flink/watermark/slow.sql 의 출력을 /root/flink/watermark/slow.out 에 저장하세요. - /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 · CREATE — WATERMARK · Time Attributes · Windowing TVF · Built-in Functions · Configuration
ts 를 이벤트 시간으로 선언한다
flink-up 으로 클러스터를 띄우고, /root/flink/watermark/ddl.sql 에 원본 CSV 를 읽는 events 표(열 event_id BIGINT, sensor STRING, reading INT, ts TIMESTAMP(3), WATERMARK FOR ts AS ts - INTERVAL '5' SECOND)와 DESCRIBE events; 를 써서 출력을 /root/flink/watermark/ddl.out 에 저장하세요.
WATERMARK 절은 열 목록 안, 마지막 열 뒤에 둡니다. DESCRIBE 결과에서 ts 의 타입 옆에 ROWTIME 이 붙고 watermark 칸에 식이 보이면 이벤트 시간 속성이 된 것입니다. filesystem 커넥터의 path 는 file:///opt/lab/fixtures/data/watermark_events.csv 입니다.
배치로 기준선을 만든다
/root/flink/watermark/batch.sql 에 SET 'execution.runtime-mode' = 'batch'; · 1단계의 events 표 · 1분 TUMBLE 창마다 window_start, window_end, COUNT(*) AS cnt, SUM(reading) AS total 을 내는 집계를 쓰고, 출력을 /root/flink/watermark/batch.out 에 저장하세요.
배치는 입력을 다 모은 뒤 계산하므로 워터마크로 행을 버리지 않습니다. 그래서 이 결과가 '늦은 행이 하나도 없었다면' 의 기준선이 됩니다. 배치에서도 TUMBLE 은 TIMESTAMP 열에 쓸 수 있습니다. cnt 를 모두 더하면 원본 행 수와 같아야 합니다.
지연 5초 스트리밍 — 버려지는 행
2단계와 같은 집계를 SET 'execution.runtime-mode' = 'streaming'; 으로 돌리는 /root/flink/watermark/w5.sql 을 만들고(워터마크는 5초 지연 그대로) 출력을 /root/flink/watermark/w5.out 에 저장하세요.
창은 window_end 가 워터마크 이하가 되는 순간 결과를 한 번 내고 상태를 비웁니다. 그 뒤에 그 창에 들어갈 행이 오면 버립니다. 결과 표의 cnt 합을 배치와 비교해 보세요. 창 결과는 +I 로만 나옵니다.
지연 0초와 30초를 한 번에 비교한다
/root/flink/watermark/sweep.sql 에 워터마크가 ts 인 표 events_0 과 ts - INTERVAL '30' SECOND 인 표 events_30(열·원천은 events 와 같음)을 만들고, 스트리밍 모드로 같은 1분 집계를 events_0 먼저, events_30 다음 순서로 돌려 출력을 /root/flink/watermark/sweep.out 에 저장하세요.
워터마크 식은 표마다 붙으므로 지연을 바꾸려면 표를 따로 만듭니다. 한 파일의 SELECT 두 개는 잡 두 개로 차례로 돌고, 출력에 결과 표가 그 순서로 찍힙니다. 지연이 0 이면 조금만 늦어도 버려지고, 30초면 대부분 기다려 줍니다.
CURRENT_WATERMARK 로 늦은 행을 뽑는다
스트리밍 모드로 5초 지연 events 에서 SELECT event_id, ts, CURRENT_WATERMARK(ts) AS wm ... WHERE CURRENT_WATERMARK(ts) IS NOT NULL AND ts <= CURRENT_WATERMARK(ts) 를 돌리는 /root/flink/watermark/late.sql 을 만들고 출력을 /root/flink/watermark/late.out 에 저장하세요.
CURRENT_WATERMARK 는 그 행이 지나는 연산자의 현재 워터마크입니다. 워터마크 생성기는 행을 먼저 내보내고 워터마크를 올리므로, 한 행이 보는 값은 그 앞 행들까지의 최대 ts − 5초입니다. 늦은 행의 수를 3단계에서 버려진 행 수와 비교해 보세요 — 늦었다고 해서 창이 닫혀 있는 것은 아닙니다.
창 앞에서 거르면 더 많이 버린다
5초 지연 events 에서 CURRENT_WATERMARK(ts) IS NULL OR ts > CURRENT_WATERMARK(ts) 인 행만 남긴 뒤(뷰나 서브쿼리) 같은 1분 TUMBLE 집계를 스트리밍으로 내는 /root/flink/watermark/filtered.sql 을 만들고 출력을 /root/flink/watermark/filtered.out 에 저장하세요.
문서가 늦은 행을 걸러 낼 때 쓰라고 권하는 식입니다. 행 단위로 '워터마크보다 이른가' 를 보므로, 창이 아직 열려 있어 받아 줬을 행까지 버립니다. 창 TVF 의 입력에 뷰를 넣으려면 CREATE VIEW 뒤 TUMBLE(TABLE 뷰, ...) 로 씁니다.
워터마크 간격을 1시간으로 — 전진이 멈춘다
3단계의 w5.sql 맨 앞에 SET 'pipeline.auto-watermark-interval' = '1 h'; 한 줄만 더한 /root/flink/watermark/slow.sql 을 돌려 출력을 /root/flink/watermark/slow.out 에 저장하세요.
워터마크는 주기(이 설정)마다 나가고, 그 사이에는 새 값이 마지막으로 낸 값보다 간격보다 더 앞설 때만 레코드에서 바로 나갑니다. 간격이 1시간이면 첫 워터마크(그전에 낸 것이 없어 바로 나감) 뒤로는 파일이 끝날 때까지 전진하지 않습니다. 그러면 닫히는 창이 있을까요?
보고서 — 지연과 버려진 행
/root/flink/watermark/report.json 에 total_rows(batch.out 의 cnt 합), dropped_0·dropped_5·dropped_30(배치 합에서 각 지연의 cnt 합을 뺀 값), late_rows_5(late.out 의 행 수), dropped_slow(배치 합 − slow.out 의 cnt 합)를 정수로 적으세요.
모두 저장한 출력에서 셀 수 있습니다. 결과 행은 '| +I |' 로 시작하고, 배치 결과 행은 날짜로 시작합니다. awk -F'|' 로 cnt 칸을 더하면 됩니다. sweep.out 에는 결과 표가 둘이니 머리줄(op 가 든 줄)을 기준으로 나눠 세세요.