Apache Flink — ストリームを本物のエンジンで動かす
総当たりログインとカード試しをパターンで捕まえる
目標
ログインと決済のイベントから、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に保存してください。- /root/flink/cep/brute.sqlに、ブルートフォースのパターン(
F{3,} S、AFTER MATCH SKIP PAST LAST ROW)のクエリを書いて実行し、出力を、/root/flink/cep/brute.outに保存してください。列はuser_id, first_fail, fails, ok_tsです。 - /root/flink/cep/nextrow.sqlに、戦略だけを
SKIP TO NEXT ROWに変えたクエリを書いて実行し、出力を、/root/flink/cep/nextrow.outに保存してください。 - /root/flink/cep/within.sqlに、ステップ2のパターンに
WITHIN INTERVAL '30' SECONDを加えたクエリを書いて実行し、出力を、/root/flink/cep/within.outに保存してください。 - /root/flink/cep/greedy.sqlに、カードテストのパターン(
S T+ B、貪欲)のクエリを書いて実行し、出力を、/root/flink/cep/greedy.outに保存してください。列はuser_id, login_ts, small_n, big_amount, big_tsです。 - /root/flink/cep/reluctant.sqlに、Tの量指定子だけを非貪欲(
T+?)に変えたクエリを書いて実行し、出力を、/root/flink/cep/reluctant.outに保存してください。 - /root/flink/cep/ddl-shuffled.sqlに
cep_events_shuffled.csvをウォーターマークの遅延10秒で読む定義を、/root/flink/cep/ddl-shuffled0.sqlに0秒で読む定義を作成し、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テーブルを書き、/root/flink/cep/count.sqlにSELECT kind, COUNT(*) AS n FROM events GROUP BY kind;を書いて、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が混ざったチェンジログで、採点ツールは、ログを最後まで適用した最終件数を、元のデータと照合します。
失敗3回以上のあとの成功
/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はその直後の行です。間に決済が入ると、マッチではありません。失敗が4回続くと、1番目・2番目の失敗から始まった候補が同じ成功で終わりますが、PAST LAST ROWは、先に始まった1つだけを出力します。
SKIP TO NEXT ROWに変える
/root/flink/cep/nextrow.sqlに、brute.sqlのAFTER MATCH SKIP PAST LAST ROWだけをAFTER MATCH SKIP TO NEXT ROWに変えたクエリを書いて実行し、出力を、/root/flink/cep/nextrow.outに保存してください。
TO NEXT ROWは、マッチの開始行の次の行から探し直します。失敗が長く続いた区間は、開始行だけが違うマッチを複数出力します。マッチ数が何倍に増えるかを、brute.outと比べてみてください。
30秒以内に終わったものだけ
/root/flink/cep/within.sqlに、brute.sqlのPATTERN (F{3,} S)の後ろにWITHIN INTERVAL '30' SECONDを付けたクエリを書いて実行し(戦略はPAST LAST ROWのまま)、出力を、/root/flink/cep/within.outに保存してください。
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で終わる行が、たいてい大きな決済であることを確認してください。
非貪欲に変える
/root/flink/cep/reluctant.sqlに、greedy.sqlのT+だけをT+?に変えたクエリを書いて実行し、出力を、/root/flink/cep/reluctant.outに保存してください。
非貪欲は、Tを最小(1つ)だけ取り込み、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に作成し、brute.sqlをそれぞれ-iで動かした出力を、/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は表の1つの列にあります。採点ツールは、元のデータから同じルールで計算し直した値と比較します。