Apache Flink — ストリームを本物のエンジンで動かす
ウィンドウ TVF — 行にウィンドウ列を3つ付ける関数
一言でいうと
Flink SQLのウィンドウは、テーブルを受け取ってテーブルを返す関数(ウィンドウTVF)です。TUMBLE・HOP・CUMULATE・SESSIONは、元の行にwindow_start・window_end・window_timeの3つの列を付けて返すだけで、集計とTop-Nは、その列でまとめる普通のSQLです。ウィンドウの種類が決めるのは、ただ1つ、1つの行が、何個の、どのウィンドウに入るかです。
なぜ必要なのか
前のモジュールで見たとおり、GROUP BY user_idのような無限の集計は、結果を際限なく書き換え、キーごとに状態を永遠に保持します。ダッシュボードが求めるのは、たいていそれではありません。「10分ごとの店別売上」「5分ごとに更新される直近10分の注文数」「今日の0時から現在までの累積売上」「1回入ってきてから出ていくまでのセッション」。どれも、時間で区切った区間の上の集計です。区間が閉じたら、結果を1回出して、状態を捨てればよいのです。
以前のFlink SQLは、これをGROUP BY TUMBLE(ts, ...)のような特殊な構文(Grouped Window Functions)で行っていました。集計には使えましたが、ウィンドウごとに順位を付けたり、2つのストリームを同じウィンドウで結合したりすることには使えませんでした。公式ドキュメント(Windowing TVF)は、ウィンドウTVFがその構文に置き換わると記しています。ウィンドウを「行に列を付ける関数」に変えたことで、ウィンドウの上に何でも載せられるようになりました。ウィンドウ集計、ウィンドウTop-N、ウィンドウ結合、ウィンドウ重複排除です。
どう動くのか
ウィンドウTVFは、FROMの位置に書きます。最初の引数はテーブル、2番目は時間属性の列、残りはサイズです。
SELECT window_start, window_end, shop, COUNT(*) AS cnt
FROM TUMBLE(TABLE orders, DESCRIPTOR(ts), INTERVAL '10' MINUTE)
GROUP BY window_start, window_end, shop;
| TVF | 引数 | 1つの行が入るウィンドウの数 | 重なり |
|---|---|---|---|
TUMBLE |
size | 1 | なし |
HOP |
slide, size | size / slide | あり |
CUMULATE |
step, size | 開始が同じウィンドウのうち、終わりがその行より後のものすべて | あり |
SESSION |
(PARTITION BYのキー)とgap | 1 | なし。サイズはまちまち |
いくつかのルールが、結果を決めます。
ウィンドウは半開区間です。[window_start, window_end)で、09:10:00ちょうどの注文は、[09:00, 09:10)ではなく[09:10, 09:20)に入ります。window_timeは、ドキュメントのとおり常にwindow_end − 1msで、ストリーミングでは、この列が時間属性として残り、次のウィンドウ演算に使えます。逆に、window_start・window_endは普通のタイムスタンプになり、時間属性ではありません。
ウィンドウの開始はエポックに揃えられます。10分ウィンドウは、正時を基準に00・10・20分に始まり、最初の行が09:01:35に来ても、ウィンドウは09:00に始まります。これをずらすのが、オプション引数のoffsetです。
HOPとCUMULATEは、引数の順序が落とし穴です。どちらも小さい値(slide、step)が先で、大きい値(size)があとです。逆に書くと、ジョブが起動する前に拒否されます。ラボ環境では、HOPは「size must be an integral multiple of slide」、CUMULATEは「maxSize must be an integral multiple of step」というエラーが出ました。つまり、sizeはslide(step)の整数倍でなければなりません。HOP(5分、10分)では、1つの行が2つのウィンドウに入るので、ウィンドウごとの件数をすべて足すと、元の行数の2倍になります。これはバグではなく、定義です。CUMULATEは、ドキュメントの表現どおり、「sizeでTUMBLEしたあと、その中をstepごとに終わりが伸びるウィンドウに分けたもの」なので、1時間のウィンドウでstepが10分なら、開始が同じウィンドウが6つ出ます。
SESSIONにはサイズがありません。同じキーの隣り合う行の間隔がgap以下なら、1つのセッションとしてつながり、セッションの終わりは、最後の行 + gapです。そのため、セッションごとに長さが違います。ドキュメントによると、SESSION TVFはまだバッチモードをサポートしていません。
ウィンドウ演算は、最後に1回だけ出力します。ドキュメント(Window Aggregation)は、ウィンドウ集計が途中の結果を出さず、ウィンドウが終わったときに最終結果だけを出し、不要になった状態をすべて消すと記しています。そのため、結果のログには+Iしかありません。ウィンドウが「終わった」という判断は、前のモジュールのウォーターマークが行います。
ウィンドウTop-Nは、ウィンドウの列で分けます。ROW_NUMBER() OVER (PARTITION BY window_start, window_end ORDER BY ...)のように、PARTITION BYにウィンドウの列2つがあってはじめて、オプティマイザーがウィンドウTop-Nに変換します。すると、通常のTop-Nと違って、順位が変わるたびに-U/+Uを出さず、ウィンドウが閉じるときに上位N件だけを1回出力します。ウィンドウ集計の上に置くことも、ウィンドウTVFの直上に置いて行そのものの順位を付けることもできます。ドキュメントによると、TVFの直上のウィンドウTop-Nは、TUMBLE・HOP・CUMULATEだけをサポートしています。
現場での姿
「5分ごとに直近1時間」のダッシュボードをHOPで作ると、1つの行が12個のウィンドウに入ります。状態と計算がそのぶん増え、ウィンドウごとの件数を合計して「総注文数」として使うと、12倍に膨らんでいます。累積指標(「今日のここまで」)をHOPで真似するケースもよくありますが、その場合は、開始が固定されたCUMULATEが正解です。
2つ目は、順位の揺れです。金額が同じ注文2つが3位を争うと、ROW_NUMBERは2つのうち1つを選びます。どちらかは保証されないので、再処理すると変わることがあります。ソートキーに注文番号のような補助キーを加えて、同点を解消しておく習慣が必要です。このラボのデータは、金額とウィンドウごとの売上に、同点がないように作ってあります。
3つ目は、タイムゾーンです。このラボはTIMESTAMP(3)なので、ウィンドウは書かれた時刻のまま切られます。TIMESTAMP_LTZの列で1日のウィンドウを作ると、セッションのタイムゾーンによって「1日」の境界が変わります。日別のウィンドウが韓国時間の午前9時で切られていたら、まずこれを疑います。
次のラボですること
注文119件に10分のTUMBLEをかけて、ウィンドウの列がどう付くかを見て、店ごとの10分集計を出します。HOP(5分、10分)の件数の合計が2倍になること、CUMULATE(10分、1時間)が開始の同じウィンドウを6つ出すこと、SESSION(gap 5分)が店ごとに長さの違うセッションを作ることを確認します。10分ウィンドウごとの売上上位2店と、30分ウィンドウごとの大きな注文3件を、ウィンドウTop-Nで取り出し、数字を報告書にまとめます。