用模式捕获暴力登录与盗卡试刷
目标
用 MATCH_RECOGNIZE 在登录和支付事件中找出暴力破解(失败 3 次以上后成功)与盗卡试刷(登录后小额支付再大额支付),并确认 AFTER MATCH SKIP 策略、WITHIN、贪婪/非贪婪以及水位线延迟如何改变匹配。
为什么重要
检测规则像“连续 · 之后 · 多长时间内”一样,是针对行顺序的条件,用 GROUP BY 或窗口聚合写不出来。MATCH_RECOGNIZE 让你能以声明式方式书写,但同一个模式也会因策略、量词和时间限制的不同,使告警数量相差数倍。本实验的评分器不会查询集群,而是对源文件按同样的规则用 Python 重新找出匹配,再与你的 sql-client 输出逐行核对。
步骤
- 运行
flink-up后,创建 /root/flink/cep/ddl.sql,其中定义读取cep_events.csv的events表(水位线为ts - INTERVAL '1' SECOND);再创建 /root/flink/cep/count.sql,按类型(kind)输出条数n,并把它的输出保存到 /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。 - 创建以水位线延迟 10 秒读取
cep_events_shuffled.csv的 /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,并分别用 -i 运行 brute.sql,把输出保存到 /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 是表中的一列。评分器会与从源数据按同样规则重新计算的值进行比较。