LabHub
배우기 러닝패스 코스

Apache Flink — Running Streams on a Real Engine

Count the Dropped Rows While Changing the Delay

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 를 적으세요.

참고

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.sqlSET '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_0ts - 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.jsontotal_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 가 든 줄)을 기준으로 나눠 세세요.