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

참고

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