Apache Flink — Running Streams on a Real Engine
Catch brute-force logins and card testing with patterns
한국어 원문으로 표시합니다.
목표
로그인·결제 사건에서 MATCH_RECOGNIZE 로 무차별 대입(실패 3번 이상 뒤 성공)과 카드 시험(로그인 뒤 소액 결제 뒤 큰 결제)을 찾고, AFTER MATCH SKIP 전략 · WITHIN · 탐욕/비탐욕 · 워터마크 지연이 매치를 어떻게 바꾸는지 확인한다.
왜 중요한가
탐지 규칙은 "연달아 · 그 뒤에 · 얼마 안에" 처럼 행의 순서에 대한 조건이라 GROUP BY 나 창 집계로는 못 쓴다. MATCH_RECOGNIZE 는 이것을 선언적으로 쓰게 해 주지만, 같은 패턴도 전략·수량자·시간 제한에 따라 경보 수가 몇 배로 달라진다. 이 실습의 채점기는 클러스터에 묻지 않는다. 원본 파일에서 같은 규칙으로 매치를 파이썬으로 다시 찾아 여러분의 sql-client 출력과 한 행씩 대조한다.
단계
flink-up뒤 /root/flink/cep/ddl.sql 에cep_events.csv를 읽는events표(워터마크ts - INTERVAL '1' SECOND)를 만들고, 종류(kind)별 건수n을 내는 /root/flink/cep/count.sql 의 출력을 /root/flink/cep/count.out 에 저장하세요.- 무차별 대입 패턴(
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. - 전략만
SKIP TO NEXT ROW로 바꾼 /root/flink/cep/nextrow.sql 의 출력을 /root/flink/cep/nextrow.out 에 저장하세요. - 2단계 패턴에
WITHIN INTERVAL '30' SECOND를 더한 /root/flink/cep/within.sql 의 출력을 /root/flink/cep/within.out 에 저장하세요. - 카드 시험 패턴(
S T+ B, 탐욕)의 /root/flink/cep/greedy.sql 출력을 /root/flink/cep/greedy.out 에 저장하세요. 열은user_id, login_ts, small_n, big_amount, big_ts. - T 수량자만 비탐욕(
T+?)으로 바꾼 /root/flink/cep/reluctant.sql 의 출력을 /root/flink/cep/reluctant.out 에 저장하세요. 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 에 저장하세요.- /root/flink/cep/report.json 에 매치 수와 금액 합을 모으세요.
참고
- 원본 열:
user_id STRING, kind STRING, amount INT, ts TIMESTAMP(3)(머리글 없는 CSV).kind는FAIL(로그인 실패) ·OK(로그인 성공) ·PAY(결제). ts 는 파일 전체에서 겹치지 않는 초 단위입니다. cep_events_shuffled.csv는 같은 행을 도착 순서만 흐트러뜨린 파일입니다. 어떤 행도 자기보다 8초 넘게 뒤의 행보다 늦게 도착하지 않습니다.- 표 정의를 되풀이하지 않으려면
sql-client.sh -i ddl.sql -f brute.sql > brute.out 2>&1. 스트리밍 결과라 출력 맨 앞에op열(+I)이 붙고 끝에Received a total of N rows가 찍힙니다. - 카드 시험 조건:
S AS S.kind = 'OK',T AS T.kind = 'PAY' AND T.amount < 20,B AS B.kind = 'PAY' AND B.amount >= 10. MEASURES 는S.ts AS login_ts, COUNT(T.amount) AS small_n, B.amount AS big_amount, B.ts AS big_ts. - 흔한 실수: 패턴 변수 사이에 다른 행이 끼면 매치가 안 됩니다(엄격한 연속). 마지막 변수에는 탐욕 수량자를 붙일 수 없습니다.
- 공식 문서: Pattern Recognition · Time Attributes · Timely Stream Processing
사건 표를 만들고 센다
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.sql 을 sql-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.sql 에 PARTITION 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 ROW 만 AFTER 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.json 에 past_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 는 표의 한 열에 있습니다. 채점기는 원본에서 같은 규칙으로 다시 계산한 값과 비교합니다.