Apache Flink — Running Streams on a Real Engine
Catch brute-force logins and card testing with patterns
Goal
From login and payment events, find brute-force attacks (a success after three or more failures) and card testing (a large payment after a small payment after a login) with MATCH_RECOGNIZE, and confirm how the AFTER MATCH SKIP strategy, WITHIN, greedy/reluctant quantifiers, and the watermark delay change the matches.
Why it matters
A detection rule is a condition on the order of rows, like "in a row · after that · within so long", so it cannot be written with GROUP BY or a window aggregation. MATCH_RECOGNIZE lets you write it declaratively, but even the same pattern can change the number of alerts severalfold depending on the strategy, the quantifier, and the time limit. The grader of this lab does not ask the cluster anything. It finds the matches again in Python from the source file with the same rules and compares them row by row with your sql-client output.
Steps
- After
flink-up, create aneventstable that readscep_events.csv(watermarkts - INTERVAL '1' SECOND) in /root/flink/cep/ddl.sql, and save the output of /root/flink/cep/count.sql, which produces the countnper kind (kind), to /root/flink/cep/count.out. - Save the output of /root/flink/cep/brute.sql for the brute-force pattern (
F{3,} S,AFTER MATCH SKIP PAST LAST ROW) to /root/flink/cep/brute.out. The columns areuser_id, first_fail, fails, ok_ts. - Save the output of /root/flink/cep/nextrow.sql, with only the strategy changed to
SKIP TO NEXT ROW, to /root/flink/cep/nextrow.out. - Save the output of /root/flink/cep/within.sql, which adds
WITHIN INTERVAL '30' SECONDto the pattern of step 2, to /root/flink/cep/within.out. - Save the output of /root/flink/cep/greedy.sql for the card-testing pattern (
S T+ B, greedy) to /root/flink/cep/greedy.out. The columns areuser_id, login_ts, small_n, big_amount, big_ts. - Save the output of /root/flink/cep/reluctant.sql, with only the T quantifier changed to reluctant (
T+?), to /root/flink/cep/reluctant.out. - Create /root/flink/cep/ddl-shuffled.sql, which reads
cep_events_shuffled.csvwith a watermark delay of 10 seconds, and /root/flink/cep/ddl-shuffled0.sql, which reads it with 0 seconds, run brute.sql with each, and save the output to /root/flink/cep/shuffled.out and /root/flink/cep/shuffled0.out. - In /root/flink/cep/report.json, collect the number of matches and the sums of amounts.
Notes
- Source columns:
user_id STRING, kind STRING, amount INT, ts TIMESTAMP(3)(a CSV with no header).kindisFAIL(login failure),OK(login success), orPAY(payment). ts is in seconds and does not overlap anywhere in the whole file. cep_events_shuffled.csvis a file with the same rows and only the arrival order scrambled. No row arrives later than a row more than 8 seconds after it.- To avoid repeating the table definition, use
sql-client.sh -i ddl.sql -f brute.sql > brute.out 2>&1. Since it is a streaming result, anopcolumn (+I) is attached at the front of the output andReceived a total of N rowsis printed at the end. - Card-testing conditions:
S AS S.kind = 'OK',T AS T.kind = 'PAY' AND T.amount < 20,B AS B.kind = 'PAY' AND B.amount >= 10. MEASURES isS.ts AS login_ts, COUNT(T.amount) AS small_n, B.amount AS big_amount, B.ts AS big_ts. - A common mistake: if another row comes between pattern variables, there is no match (strict contiguity). You cannot attach a greedy quantifier to the last variable.
- Official docs: Pattern Recognition · Time Attributes · Timely Stream Processing
Create the event table and count
After flink-up, in /root/flink/cep/ddl.sql write an events table that reads /opt/lab/fixtures/data/cep_events.csv and has WATERMARK FOR ts AS ts - INTERVAL '1' SECOND, and create /root/flink/cep/count.out by running /root/flink/cep/count.sql, which contains SELECT kind, COUNT(*) AS n FROM events GROUP BY kind;, with sql-client.sh -i ddl.sql -f count.sql > count.out 2>&1.
The ORDER BY of MATCH_RECOGNIZE must be a time attribute, so put the watermark on ts. Since it is a streaming GROUP BY, the output is a changelog with -U/+U mixed in, and the grader applies the log to the end and compares the final counts with the source.
A success after three or more failures
In /root/flink/cep/brute.sql, write a query with 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), and DEFINE F AS F.kind = 'FAIL', S AS S.kind = 'OK', and save the output to /root/flink/cep/brute.out.
F{3,} is three or more, and S is the very next row. If a payment comes in between, it is not a match. If failures continue four times, the candidate starting at the first failure and the candidate starting at the second failure end at the same success, but PAST LAST ROW emits only the one that started first.
Switch to SKIP TO NEXT ROW
Run /root/flink/cep/nextrow.sql, which is brute.sql with only AFTER MATCH SKIP PAST LAST ROW changed to AFTER MATCH SKIP TO NEXT ROW, and save it to /root/flink/cep/nextrow.out.
TO NEXT ROW searches again from the row after the start row of the match. A stretch where failures continue for long emits several matches that differ only in the start row. Compare with brute.out how many times the number of matches grows.
Only those that finished within 30 seconds
Run /root/flink/cep/within.sql, which is brute.sql with WITHIN INTERVAL '30' SECOND attached after PATTERN (F{3,} S), and save it to /root/flink/cep/within.out (the strategy stays PAST LAST ROW).
WITHIN limits the interval between the first row and the last row of the match. A slow attempt drops out entirely, and in a stretch where failures continue for long, a match may remain with its start row pushed back because the leading failures were cut off. A candidate whose interval is exactly 30 seconds drops out in this engine.
Card testing — a greedy quantifier
In /root/flink/cep/greedy.sql, write a query with the conditions and MEASURES from the Notes and PATTERN (S T+ B) and AFTER MATCH SKIP PAST LAST ROW, and save the output to /root/flink/cep/greedy.out. The columns are user_id, login_ts, small_n, big_amount, big_ts.
T (a payment under 20) and B (a payment of 10 or more) overlap at 10–19. The default quantifier is greedy, so it eats T as far as it can, and the next row must be a B for it to be a match. Confirm that the row ending in B is usually a large payment.
Switch to reluctant
Run /root/flink/cep/reluctant.sql, which is greedy.sql with only T+ changed to T+?, and save it to /root/flink/cep/reluctant.out.
Reluctant eats T only the minimum (one) and ends the match at the first row that can be a B. Payments of 10–19 are now caught as B, so the match ends earlier and big_amount gets smaller. Compare the number of matches and the big_amount sum with greedy.out.
Scrambled arrival order and the watermark
Based on ddl.sql, create /root/flink/cep/ddl-shuffled.sql with the path changed to cep_events_shuffled.csv and the watermark delay changed to 10 seconds, and /root/flink/cep/ddl-shuffled0.sql with it changed to 0 seconds, run brute.sql with each using -i, and save the output to /root/flink/cep/shuffled.out and /root/flink/cep/shuffled0.out.
Event-time MATCH_RECOGNIZE sorts the rows and then searches the pattern, but the sort waits only as far as the watermark. If the delay is larger than the scramble (at most 8 seconds), the result is the same as the sorted file, and if the delay is 0, rows whose event time is earlier than a row that arrived before them are dropped as late rows, so matches decrease.
Report — the rules decide the number of alerts
In /root/flink/cep/report.json, write as integers past_last_row (the number of matches in brute.out), to_next_row (nextrow.out), within_30s (within.out), greedy_big_total and reluctant_big_total (the sum of big_amount in each output), and matches_lost_without_delay (the number of matches in shuffled.out − the number of matches in shuffled0.out).
The number of matches is in Received a total of N rows at the end of each output, and big_amount is in one column of the table. The grader compares with the values recomputed from the source with the same rules.