LabHub
배우기 러닝패스 코스

Apache Flink — ストリームを本物のエンジンで動かす

総当たりログインとカード試しをパターンで捕まえる

LabHub 에서 이어서 보기

한국어 원문으로 표시합니다.

목표

로그인·결제 사건에서 MATCH_RECOGNIZE 로 무차별 대입(실패 3번 이상 뒤 성공)과 카드 시험(로그인 뒤 소액 결제 뒤 큰 결제)을 찾고, AFTER MATCH SKIP 전략 · WITHIN · 탐욕/비탐욕 · 워터마크 지연이 매치를 어떻게 바꾸는지 확인한다.

왜 중요한가

탐지 규칙은 "연달아 · 그 뒤에 · 얼마 안에" 처럼 행의 순서에 대한 조건이라 GROUP BY 나 창 집계로는 못 쓴다. MATCH_RECOGNIZE 는 이것을 선언적으로 쓰게 해 주지만, 같은 패턴도 전략·수량자·시간 제한에 따라 경보 수가 몇 배로 달라진다. 이 실습의 채점기는 클러스터에 묻지 않는다. 원본 파일에서 같은 규칙으로 매치를 파이썬으로 다시 찾아 여러분의 sql-client 출력과 한 행씩 대조한다.

단계

  1. flink-up/root/flink/cep/ddl.sqlcep_events.csv 를 읽는 events 표(워터마크 ts - INTERVAL '1' SECOND)를 만들고, 종류(kind)별 건수 n 을 내는 /root/flink/cep/count.sql 의 출력을 /root/flink/cep/count.out 에 저장하세요.
  2. 무차별 대입 패턴(F{3,} S, AFTER MATCH SKIP PAST LAST ROW)의 /root/flink/cep/brute.sql 출력을 /root/flink/cep/brute.out 에 저장하세요. 열은 user_id, first_fail, fails, ok_ts.
  3. 전략만 SKIP TO NEXT ROW 로 바꾼 /root/flink/cep/nextrow.sql 의 출력을 /root/flink/cep/nextrow.out 에 저장하세요.
  4. 2단계 패턴에 WITHIN INTERVAL '30' SECOND 를 더한 /root/flink/cep/within.sql 의 출력을 /root/flink/cep/within.out 에 저장하세요.
  5. 카드 시험 패턴(S T+ B, 탐욕)의 /root/flink/cep/greedy.sql 출력을 /root/flink/cep/greedy.out 에 저장하세요. 열은 user_id, login_ts, small_n, big_amount, big_ts.
  6. T 수량자만 비탐욕(T+?)으로 바꾼 /root/flink/cep/reluctant.sql 의 출력을 /root/flink/cep/reluctant.out 에 저장하세요.
  7. cep_events_shuffled.csv 를 워터마크 지연 10초로 읽는 /root/flink/cep/ddl-shuffled.sql, 0초로 읽는 /root/flink/cep/ddl-shuffled0.sql 을 만들어 brute.sql 을 각각 돌린 출력을 /root/flink/cep/shuffled.out·/root/flink/cep/shuffled0.out 에 저장하세요.
  8. /root/flink/cep/report.json 에 매치 수와 금액 합을 모으세요.

참고

사건 표를 만들고 센다

flink-up/root/flink/cep/ddl.sql/opt/lab/fixtures/data/cep_events.csv 를 읽고 WATERMARK FOR ts AS ts - INTERVAL '1' SECOND 를 둔 events 표를 쓰고, SELECT kind, COUNT(*) AS n FROM events GROUP BY kind; 를 담은 /root/flink/cep/count.sqlsql-client.sh -i ddl.sql -f count.sql > count.out 2>&1 로 돌려 /root/flink/cep/count.out 을 만드세요.

MATCH_RECOGNIZE 의 ORDER BY 는 시간 속성이어야 하므로 ts 에 워터마크를 둡니다. 스트리밍 GROUP BY 라 출력은 -U/+U 가 섞인 변경 로그이고, 채점기는 로그를 끝까지 적용한 최종 건수를 원본과 대조합니다.

실패 세 번 이상 뒤 성공

/root/flink/cep/brute.sqlPARTITION BY user_id ORDER BY ts, MEASURES FIRST(F.ts) AS first_fail, COUNT(F.ts) AS fails, S.ts AS ok_ts, ONE ROW PER MATCH, AFTER MATCH SKIP PAST LAST ROW, PATTERN (F{3,} S), DEFINE F AS F.kind = 'FAIL', S AS S.kind = 'OK' 인 질의를 쓰고 출력을 /root/flink/cep/brute.out 에 저장하세요.

F{3,} 는 3번 이상, S 는 그 바로 다음 행입니다. 사이에 결제가 끼면 매치가 아닙니다. 실패가 네 번 이어지면 첫째·둘째 실패에서 시작한 후보가 같은 성공에서 끝나는데, PAST LAST ROW 는 먼저 시작한 하나만 냅니다.

SKIP TO NEXT ROW 로 바꾼다

brute.sql 에서 AFTER MATCH SKIP PAST LAST ROWAFTER MATCH SKIP TO NEXT ROW 로 바꾼 /root/flink/cep/nextrow.sql 을 돌려 /root/flink/cep/nextrow.out 에 저장하세요.

TO NEXT ROW 는 매치의 시작 행 다음 행에서 다시 찾습니다. 실패가 길게 이어진 구간은 시작 행만 다른 매치를 여러 개 냅니다. 매치 수가 몇 배로 느는지 brute.out 과 비교해 보세요.

30초 안에 끝난 것만

brute.sql 의 PATTERN (F{3,} S) 뒤에 WITHIN INTERVAL '30' SECOND 를 붙인 /root/flink/cep/within.sql 을 돌려 /root/flink/cep/within.out 에 저장하세요(전략은 PAST LAST ROW 그대로).

WITHIN 은 매치의 첫 행과 마지막 행의 간격을 제한합니다. 느린 시도는 통째로 떨어지고, 실패가 길게 이어진 구간에서는 앞머리 실패가 잘려 시작 행이 뒤로 밀린 매치가 남기도 합니다. 간격이 정확히 30초인 후보는 이 엔진에서 떨어집니다.

카드 시험 — 탐욕 수량자

/root/flink/cep/greedy.sql 에 참고의 조건과 MEASURES 로 PATTERN (S T+ B) · AFTER MATCH SKIP PAST LAST ROW 인 질의를 쓰고 출력을 /root/flink/cep/greedy.out 에 저장하세요. 열은 user_id, login_ts, small_n, big_amount, big_ts 입니다.

T(20 미만 결제)와 B(10 이상 결제)는 10–19 에서 겹칩니다. 기본 수량자는 탐욕이라 T 를 될 수 있는 데까지 먹고, 그다음 행이 B 여야 매치가 됩니다. B 로 끝나는 행이 대개 큰 결제라는 것을 확인하세요.

비탐욕으로 바꾼다

greedy.sql 에서 T+T+? 로 바꾼 /root/flink/cep/reluctant.sql 을 돌려 /root/flink/cep/reluctant.out 에 저장하세요.

비탐욕은 T 를 최소(하나)만 먹고, B 가 될 수 있는 첫 행에서 매치를 끝냅니다. 10–19 결제가 이제 B 로 잡혀 매치가 일찍 끝나고 big_amount 가 작아집니다. 매치 수와 big_amount 합을 greedy.out 과 비교하세요.

흐트러진 도착 순서와 워터마크

ddl.sql 을 바탕으로 경로를 cep_events_shuffled.csv 로, 워터마크 지연을 10초로 바꾼 /root/flink/cep/ddl-shuffled.sql 과 0초로 바꾼 /root/flink/cep/ddl-shuffled0.sql 을 만들고, brute.sql 을 각각 -i 로 돌린 출력을 /root/flink/cep/shuffled.out·/root/flink/cep/shuffled0.out 에 저장하세요.

이벤트 시간 MATCH_RECOGNIZE 는 행을 정렬한 뒤 패턴을 찾지만, 정렬은 워터마크까지만 기다립니다. 흐트러짐(최대 8초)보다 지연이 크면 정렬된 파일과 결과가 같고, 지연이 0이면 앞서 온 행보다 이벤트 시간이 이른 행이 늦은 행으로 버려져 매치가 줄어듭니다.

보고서 — 규칙이 경보 수를 정한다

/root/flink/cep/report.jsonpast_last_row(brute.out 매치 수), to_next_row(nextrow.out), within_30s(within.out), greedy_big_total·reluctant_big_total(각 출력의 big_amount 합), matches_lost_without_delay(shuffled.out 매치 수 − shuffled0.out 매치 수)를 정수로 적으세요.

매치 수는 각 출력 끝의 Received a total of N rows 에, big_amount 는 표의 한 열에 있습니다. 채점기는 원본에서 같은 규칙으로 다시 계산한 값과 비교합니다.