Apache Flink — ストリームを本物のエンジンで動かす
MATCH_RECOGNIZE — 複数行にまたがるパターンを SQL で探す
一言でいうと
MATCH_RECOGNIZEは、キーごとにイベント時間の順に並べた行に、正規表現のようなパターンを当てはめ、マッチした区間ごとに要約の行を1つ出力する句です。何が1つのマッチになるかは、パターンそのものよりも、量指定子の貪欲さ、AFTER MATCH SKIP戦略、WITHIN時間制限が決めます。
なぜ必要なのか
「同じユーザーがログインに3回以上続けて失敗したあとで、成功した」を探すとしましょう。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: マッチ1つが出力する列です。
FIRST・LASTで、変数に対応した行を前後から指し、COUNT・SUMのような集計も使えます。出力は、PARTITION BYの列 + MEASURESの列です。現在は、ONE ROW PER MATCHだけをサポートしています。
貪欲と非貪欲の違いは、次のとおりです。量指定子はデフォルトが貪欲(できるだけ多く)で、後ろに?を付けると非貪欲(できるだけ少なく)です。2つの変数の条件が重なるときだけ、違いが出ます。このラボのカードテストのパターンS T+ Bで、Tは「20未満の決済」、Bは「10以上の決済」なので、10–19が両方に当てはまります。貪欲なT+は、20未満の決済を最後まで取り込み、次の行がBでなければならず、非貪欲なT+?は、1つを取り込んだあと、Bになれる最初の行で終わらせます。このファイルでは、マッチ数は同じなのに、終わる行と大きな決済の金額が変わりました。
AFTER MATCH SKIPは、マッチを1つ見つけたあと、どこから探し直すかです。SKIP PAST LAST ROWは、マッチの最後の行の次から。1つの行は、多くても1つのマッチにしか入りません。SKIP TO NEXT ROWは、マッチの開始行の次から。失敗が5回続いたあとで成功すると、開始行だけが違うマッチが3つ出ます。SKIP TO LAST 변수・SKIP TO FIRST 변수(プレースホルダーは変数名です)もあります。
WITHINは、最初の行と最後の行の間隔の制限です(標準の外にあるFlinkの拡張)。超えた候補は捨てられ、ドキュメントのとおり、状態を空にする根拠になるので、ストリームではほぼ常に付けます。このPodで測ってみたところ、間隔がちょうど30秒の候補も、WITHIN INTERVAL '30' SECONDで落ちました。境界は「未満」と考えるのが安全です。もう1つあります。MATCH_RECOGNIZEは、table.exec.state.ttlに従いません。状態を減らすつまみは、WITHINです。
現場での姿
検知ルールは、ほぼ常に「何回以上・どれだけの時間内に・そのあとに」という形なので、MATCH_RECOGNIZEがよく合います。ただし、アラートの件数が、ルールよりもAFTER MATCH戦略に左右されることを知らないと、困ります。同じパターンが、PAST LAST ROWならアラート1件、TO NEXT ROWなら3件を出します。アラートシステムが「イベント数」を数えるのか「マッチ数」を数えるのかを、まず合わせる必要があります。
2つ目は、遅延行です。パターンは、イベント時間で並べ替えてから探すので、入力が乱れても大丈夫に見えますが、並べ替えはウォーターマークまでしか待ちません。ウォーターマークの遅延を0にして、乱れた入力を入れると、遅延行が静かに抜けて、マッチが減ります。エラーも警告もありません。このラボで、同じ行を順序だけ乱したファイルで、遅延10秒と0秒を比べます。
3つ目は、終わらないパターンです。条件のない変数に上限のない量指定子をかけると、すべての行がその変数に当てはまって、マッチが終わらず、状態だけが積み上がります。ドキュメントの対処法は、後ろの変数の条件を否定して入れるか、非貪欲に変えるか、WITHINをかけることです。
次のラボですること
ログインと決済のイベントテーブルを作って、種類別の件数で確認したあと、ブルートフォースのパターンを、PAST LAST ROW・TO NEXT ROW・WITHIN 30秒に変えながら動かします。カードテストのパターンを、貪欲・非貪欲で動かして、終わる行が変わることを見て、到着順を乱したファイルを、ウォーターマークの遅延10秒と0秒で読んで、マッチ数を比べたあと、すべての数字を報告書にまとめます。採点ツールは、同じルールをPythonでもう一度動かして、1行ずつ照合します。