MATCH_RECOGNIZE — 用 SQL 查找跨多行的模式
一句话总结
MATCH_RECOGNIZE 是这样一个子句:对每个键下按事件时间顺序排好的行套用类似正则表达式的模式,每找到一段匹配,就输出一行汇总。一次匹配究竟由什么构成,取决于量词是否贪婪、AFTER MATCH SKIP 策略和 WITHIN 时间限制,而不是模式本身。
为什么需要它
假设要找出“同一用户连续登录失败三次以上,随后登录成功”。GROUP BY 能统计失败次数,却无法表达“连续”和“紧接着”。窗口聚合会把时间切成一个个格子,跨越格子边界的尝试就被漏掉了。多次自连接虽然能表达,但在流处理中需要永远记住两侧的数据(第 5 个模块),查询也会变得难以阅读。
真正需要的是“针对行顺序的条件”。Flink 本来就有复杂事件处理(CEP)库,又把 SQL 标准(ISO/IEC TR 19075-5:2016)中的行模式识别建立在它之上,以 MATCH_RECOGNIZE 的形式提供(官方文档中的 Pattern Recognition)。
工作原理
骨架如下。
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 —— 按键分别查找。文档警告说,如果不分区,为了保持全局顺序,该算子的并行度会变成 1。
ORDER BY的第一列必须是升序的时间属性。如果是事件时间,行会先被排序再送入模式机,因此即使到达顺序被打乱,结果也遵循事件时间顺序——但比水位线晚到的行无法参与排序,会被丢弃。 - PATTERN · DEFINE —— 像正则表达式一样用模式变量和量词(
*+?{n}{n,}{n,m})书写,并为每个变量指定条件。相连书写的变量之间不能插入其他行(严格连续)。可能产生空匹配的模式(A*)以及以贪婪量词结尾的最后一个变量(A B*)都不被允许。 - MEASURES —— 每次匹配输出的列。用
FIRST、LAST从前后两端指向变量所匹配的行,也可以使用COUNT、SUM之类的聚合。输出是PARTITION BY的列加上 MEASURES 的列。目前只支持ONE ROW PER MATCH。
贪婪与非贪婪。 量词默认是贪婪的(尽量多匹配),在后面加 ? 就变成非贪婪(尽量少匹配)。只有当两个变量的条件重叠时,才会有差别。本实验的盗卡试刷模式 S T+ B 中,T 是“小于 20 的支付”,B 是“大于等于 10 的支付”,所以 10–19 同时满足两者。贪婪的 T+ 会一直吃到小于 20 的支付结束,下一行必须是 B;非贪婪的 T+? 吃掉一个之后,就在第一个可以成为 B 的行结束。在这个文件里,匹配数相同,但结束的行和大额支付金额都变了。
AFTER MATCH SKIP。 决定找到一次匹配后从哪里重新查找。SKIP PAST LAST ROW 从匹配的最后一行之后开始——一行最多只属于一次匹配。SKIP TO NEXT ROW 从匹配的起始行之后开始——连续失败五次后成功,会得到三个仅起始行不同的匹配。另外还有 SKIP TO LAST 변수 和 SKIP TO FIRST 변수(占位符均为变量名)。
WITHIN。 限制第一行与最后一行之间的间隔(这是标准之外的 Flink 扩展)。超出限制的候选会被丢弃,而且正如文档所说,它是清理状态的依据,因此在流处理中几乎总要加上。在这个 Pod 上实测,间隔恰好为 30 秒的候选也被 WITHIN INTERVAL '30' SECOND 淘汰了——边界最好按“小于”来理解。还有一点——MATCH_RECOGNIZE 不遵循 table.exec.state.ttl。缩减状态的旋钮是 WITHIN。
在现场相遇的样子
检测规则几乎总是“至少几次 · 在多长时间内 · 之后”的形态,所以 MATCH_RECOGNIZE 非常合适。但如果不知道告警数量取决于 AFTER MATCH 策略而不是规则本身,就会遇到麻烦。同一个模式,用 PAST LAST ROW 产生 1 条告警,用 TO NEXT ROW 则产生 3 条。首先要对齐的是:告警系统统计的是“事件数”还是“匹配数”。
第二个是迟到的行。模式是按事件时间排序之后才查找的,所以输入乱序看起来没关系,但排序只会等到水位线为止。如果把水位线延迟设为 0 并送入乱序输入,迟到的行会悄悄掉队,匹配随之减少。既没有错误,也没有警告。本实验会用同样的行、只打乱顺序的文件,对比延迟 10 秒和 0 秒的结果。
第三个是永远结束不了的模式。给没有条件的变量加上没有上限的量词,所有行都会被它匹配,匹配永远结束不了,状态只会不断堆积。文档给出的办法是:把后一个变量的条件取反写进去,或者改成非贪婪,或者加上 WITHIN。
下一项实验要做什么
创建登录和支付事件表,按类型统计条数确认之后,把暴力破解模式依次换成 PAST LAST ROW · TO NEXT ROW · WITHIN 30 秒来运行。把盗卡试刷模式分别用贪婪和非贪婪运行,观察结束的行有何不同;再用水位线延迟 10 秒和 0 秒读取到达顺序被打乱的文件,比较匹配数;最后把所有数字汇总成报告。评分器会用 Python 重新执行同样的规则,与输出逐行核对。