Apache Flink — Running Streams on a Real Engine
MATCH_RECOGNIZE — finding patterns across rows with SQL
In one line
MATCH_RECOGNIZE is a clause that tests a regular-expression-like pattern against rows lined up in event-time order for each key, and emits one summary row for each matching stretch. What becomes one match is decided not so much by the pattern itself as by whether the quantifier is greedy, the AFTER MATCH SKIP strategy, and the WITHIN time limit.
Why this was needed
Suppose you want to find "the same user failed to log in more than three times in a row and then succeeded". With GROUP BY you can count the failures, but you cannot express "in a row" and "right after that". A window aggregation cuts time into slots, so it misses attempts that straddle a slot boundary. If you apply a self-join several times it can be expressed, but in a stream you would have to remember both sides forever (module 5) and the query becomes unreadable.
What you need is "a condition on the order of rows". Flink already had a complex event processing (CEP) library, and it put row pattern recognition, which entered the SQL standard (ISO/IEC TR 19075-5:2016), on top of it and exposed it as MATCH_RECOGNIZE (Pattern Recognition in the official docs).
How it works
The skeleton looks like this.
SELECT * FROM events MATCH_RECOGNIZE (
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) WITHIN INTERVAL '30' SECOND
DEFINE F AS F.kind = 'FAIL', S AS S.kind = 'OK'
) AS T;
- PARTITION BY · ORDER BY — It searches separately for each key. The docs warn that if you do not partition, it becomes an operator with parallelism 1 to keep a global order. The first column of
ORDER BYmust be an ascending time attribute. With event time, rows are first sorted and then fed into the pattern machine, so even if the arrival order is scrambled, the result follows event-time order — but a row that arrives later than the watermark cannot join the sort and is dropped. - PATTERN · DEFINE — You write it like a regular expression with pattern variables and quantifiers (
*+?{n}{n,}{n,m}), and give each variable a condition. No other row may come between variables written one after another (strict contiguity). A pattern that can produce an empty match (A*) and a greedy quantifier on the last variable (A B*) are not allowed. - MEASURES — The columns a match emits.
FIRSTandLASTpoint to the rows bound to a variable from the front and the back, and aggregates such asCOUNTandSUMcan also be used. The output is thePARTITION BYcolumns + the MEASURES columns. For now onlyONE ROW PER MATCHis supported.
Greedy and reluctant. A quantifier is greedy by default (as many as possible), and if you add ? after it, it is reluctant (as few as possible). It makes a difference only when the conditions of two variables overlap. In the card-testing pattern S T+ B of this lab, T is "a payment under 20" and B is "a payment of 10 or more", so 10–19 fits both. A greedy T+ eats payments under 20 to the end and the next row must be a B, while a reluctant T+? eats one and ends at the first row that can be a B. In this file the number of matches was the same, but the row where it ends and the large payment amount differed.
AFTER MATCH SKIP. This is where to search again after finding one match. SKIP PAST LAST ROW searches from after the last row of the match — a row is in at most one match. SKIP TO NEXT ROW searches from after the start row of the match — if a success follows five failures in a row, you get three matches that differ only in the start row. There are also SKIP TO LAST 변수 and SKIP TO FIRST 변수 (where the placeholder stands for a pattern variable).
WITHIN. This is a limit on the interval between the first row and the last row (a Flink extension outside the standard). Candidates that exceed it are discarded, and as the docs say, it is the basis for clearing state, so in a stream you attach it almost always. When I measured on this Pod, even a candidate whose interval was exactly 30 seconds dropped out with WITHIN INTERVAL '30' SECOND — it is safest to think of the boundary as "less than". One more thing — MATCH_RECOGNIZE does not follow table.exec.state.ttl. The knob for reducing state is WITHIN.
What it looks like in the field
Detection rules almost always have the shape "at least so many times, within so long, and after that", so MATCH_RECOGNIZE fits well. But it becomes a problem if you do not know that the number of alerts depends on the AFTER MATCH strategy more than on the rule. The same pattern gives 1 alert with PAST LAST ROW and 3 alerts with TO NEXT ROW. You first have to settle whether the alerting system counts "incidents" or "matches".
The second is late rows. A pattern is searched after sorting by event time, so it seems fine even if the input is scrambled, but the sort waits only as far as the watermark. If you set the watermark delay to 0 and feed scrambled input, late rows silently drop out and matches decrease. No error, no warning. In this lab you compare delays of 10 seconds and 0 seconds on a file with the same rows but scrambled order.
The third is a pattern that never ends. If you put an unbounded quantifier on a variable with no condition, every row fits that variable, so the match never ends and only state piles up. The docs' prescription is to put in the negation of the next variable's condition, to switch to reluctant, or to attach WITHIN.
What you will do in the next lab
You create a table of login and payment events, check it with counts by kind, and then run the brute-force pattern while switching it between PAST LAST ROW, TO NEXT ROW, and WITHIN 30 seconds. You run the card-testing pattern greedy and reluctant to see the row where it ends change, read a file with scrambled arrival order with watermark delays of 10 seconds and 0 seconds to compare the number of matches, and then collect all the numbers into a report. The grader reruns the same rules in Python and compares row by row.