TT Lab
开始
学习 学习路径 课程

Apache Flink — 用真正的引擎跑流处理

用模式捕获暴力登录与盗卡试刷

在 TT Lab 中继续学习

目标

用 MATCH_RECOGNIZE 在登录和支付事件中找出暴力破解(失败 3 次以上后成功)与盗卡试刷(登录后小额支付再大额支付),并确认 AFTER MATCH SKIP 策略、WITHIN、贪婪/非贪婪以及水位线延迟如何改变匹配。

为什么重要

检测规则像“连续 · 之后 · 多长时间内”一样,是针对行顺序的条件,用 GROUP BY 或窗口聚合写不出来。MATCH_RECOGNIZE 让你能以声明式方式书写,但同一个模式也会因策略、量词和时间限制的不同,使告警数量相差数倍。本实验的评分器不会查询集群,而是对源文件按同样的规则用 Python 重新找出匹配,再与你的 sql-client 输出逐行核对。

步骤

  1. 运行 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。
  2. 把暴力破解模式(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。
  3. 只把策略改成 SKIP TO NEXT ROW,并把 /root/flink/cep/nextrow.sql 的输出保存到 /root/flink/cep/nextrow.out。
  4. 在第 2 步的模式上加入 WITHIN INTERVAL '30' SECOND,并把 /root/flink/cep/within.sql 的输出保存到 /root/flink/cep/within.out。
  5. 把盗卡试刷模式(S T+ B,贪婪)的 /root/flink/cep/greedy.sql 的输出保存到 /root/flink/cep/greedy.out。列为 user_id, login_ts, small_n, big_amount, big_ts。
  6. 只把 T 量词改成非贪婪(T+?),并把 /root/flink/cep/reluctant.sql 的输出保存到 /root/flink/cep/reluctant.out。
  7. 创建以水位线延迟 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。
  8. 把匹配数和金额合计汇总到 /root/flink/cep/report.json。

参考

创建事件表并统计条数

运行 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 是表中的一列。评分器会与从源数据按同样规则重新计算的值进行比较。