Apache Flink — ストリームを本物のエンジンで動かす
4種類のウィンドウ TVF で同じ注文を切り分ける
目標
同じ注文ストリームにTUMBLE・HOP・CUMULATE・SESSIONウィンドウをかけて、ウィンドウの列がどう付き、1つの行が何個のウィンドウに入るかを、結果で確認します。ウィンドウ集計の上のウィンドウTop-Nと、ウィンドウTVFの直上のウィンドウTop-Nを作ります。
なぜ重要なのか
ウィンドウの種類を間違えて選ぶと、数字が静かに間違います。HOPの件数を合計すると重なった分だけ膨らみ、GROUP BYからウィンドウの列を抜くと、ウィンドウ集計ではなく無限の集計になって、更新ログが混ざります。ウィンドウTVFは、行にウィンドウの列を付ける関数にすぎないので、その列がどう付くかさえ正確に知っていれば、集計でも順位でも、普通のSQLで載せられます。このラボの採点ツールは、クラスターに問い合わせません。皆さんが保存したsql-clientの出力を読み、ウィンドウごとの期待値を、元のCSVからPythonで直接計算して照合します。
ステップ
flink-upでクラスターを起動し、/root/flink/windows/assign.sqlに、10分のTUMBLEが付けたorder_id, ts, window_start, window_end, window_timeを、order_id <= 20の注文だけ取り出すクエリを書いて実行し、出力を、/root/flink/windows/assign.outに保存してください。- /root/flink/windows/tumble.sqlに、10分の
TUMBLEウィンドウ・店(shop)ごとにcnt(件数)・revenue(amountの合計)を出すクエリを書いて実行し、出力を、/root/flink/windows/tumble.outに保存してください。 - /root/flink/windows/hop.sqlに、
HOP(slide 5分、size 10分)のウィンドウごとにcntを出すクエリを書いて実行し、出力を、/root/flink/windows/hop.outに保存してください。 - /root/flink/windows/cumulate.sqlに、
CUMULATE(step 10分、size 1時間)のウィンドウごとにrevenueを出すクエリを書いて実行し、出力を、/root/flink/windows/cumulate.outに保存してください。 - /root/flink/windows/session.sqlに、
SESSION(PARTITION BY shop、gap 5分)のウィンドウ・店ごとにcntを出すクエリを書いて実行し、出力を、/root/flink/windows/session.outに保存してください。 - /root/flink/windows/top-shops.sqlに、10分のTUMBLEウィンドウごとに、店別の売上の上位2店(
window_start, window_end, shop, revenue, rownum)を出すクエリを書いて実行し、出力を、/root/flink/windows/top-shops.outに保存してください。 - /root/flink/windows/top-orders.sqlに、30分のTUMBLEウィンドウごとに、金額が大きい注文3件(
order_id, shop, amount, window_start, window_end, rownum)を集計なしで出すクエリを書いて実行し、出力を、/root/flink/windows/top-orders.outに保存してください。 - /root/flink/windows/report.jsonに、
orders・hop_assignments・cumulate_windows・sessions・max_session_ordersを書いてください。
参考
- 元データ:
/opt/lab/fixtures/data/windows_orders.csv、列はorder_id BIGINT, shop STRING, amount INT, ts TIMESTAMP(3)(ヘッダーなしのCSV、tsは昇順)。テーブルにはWATERMARK FOR ts AS ts - INTERVAL '1' SECONDを置き、ストリーミングモードで動かします。 - 形:
FROM TUMBLE(TABLE orders, DESCRIPTOR(ts), INTERVAL '10' MINUTE)。HOPは(TABLE, DESCRIPTOR, slide, size)、CUMULATEは(TABLE, DESCRIPTOR, step, size)、SESSIONは(TABLE orders PARTITION BY shop, DESCRIPTOR(ts), gap)です。 - ウィンドウ集計は、
GROUP BY window_start, window_end, ...でまとめます。ウィンドウの列を抜くと、無限の集計になって-U/+Uが混ざります。 - よくある間違い: HOP・CUMULATEの引数の順序(小さい値が先)を入れ替えること。sizeがslide(step)の整数倍ではないというエラーで拒否されます。
- よくある間違い: ウィンドウTop-NのPARTITION BYからwindow_start, window_endを抜くこと。そうすると、ウィンドウTop-Nではなく通常のTop-Nになり、順位が変わるたびにログが出ます。
- 公式ドキュメント: Windowing TVF・Window Aggregation・Window Top-N・Time Attributes
ウィンドウTVFが付ける3つの列
flink-upでクラスターを起動し、/root/flink/windows/assign.sqlに、元のテーブルorders(参考の列とウォーターマーク)とSELECT order_id, ts, window_start, window_end, window_time FROM TUMBLE(TABLE orders, DESCRIPTOR(ts), INTERVAL '10' MINUTE) WHERE order_id <= 20を書いてストリーミングで実行し、出力を、/root/flink/windows/assign.outに保存してください。
ウィンドウTVFは、元の列をそのままにして、ウィンドウの列3つを付けて返します。ウィンドウは[開始、終わり)の半開区間なので、境界の時刻ちょうどの注文は、その時刻に始まるウィンドウに入ります。window_timeとwindow_endの差を見てください。
店ごとの10分集計
/root/flink/windows/tumble.sqlに、10分のTUMBLEウィンドウとshopでまとめてwindow_start, window_end, shop, COUNT(*) AS cnt, SUM(amount) AS revenueを出すウィンドウ集計を書き、出力を、/root/flink/windows/tumble.outに保存してください。
ウィンドウ集計は、GROUP BYにwindow_startとwindow_endを入れます。ウィンドウが重ならないので、cntをすべて足すと、元の注文数と同じになります。結果は、ウィンドウが閉じるときに+Iで1回だけ出ます。
HOP: 1つの注文が2つのウィンドウに入る
/root/flink/windows/hop.sqlに、HOP(TABLE orders, DESCRIPTOR(ts), INTERVAL '5' MINUTE, INTERVAL '10' MINUTE)のウィンドウごとにwindow_start, window_end, COUNT(*) AS cntを出す集計を書き、出力を、/root/flink/windows/hop.outに保存してください。
HOPの3番目の引数がslide(ウィンドウが始まる間隔)、4番目がsize(ウィンドウの長さ)です。5分ごとに始まる10分のウィンドウなら、ウィンドウが半分ずつ重なって、1つの注文が2つのウィンドウに入ります。cntの合計を、元の注文数と比べてみてください。最初のウィンドウは、最初の注文より5分早い時刻に始まることがあります。
CUMULATE: 開始が固定された累積ウィンドウ
/root/flink/windows/cumulate.sqlに、CUMULATE(TABLE orders, DESCRIPTOR(ts), INTERVAL '10' MINUTE, INTERVAL '1' HOUR)のウィンドウごとにwindow_start, window_end, SUM(amount) AS revenueを出す集計を書き、出力を、/root/flink/windows/cumulate.outに保存してください。
CUMULATEは、size(1時間)でTUMBLEしたあと、その中をstep(10分)ごとに終わりが伸びるウィンドウに分けたものです。開始が同じウィンドウのrevenueは、終わりが遅いほど大きくなるか、同じになるはずです。1時間ごとにウィンドウがいくつ出るかを数えてみてください。
SESSION: 店ごとに長さが違うウィンドウ
/root/flink/windows/session.sqlに、SESSION(TABLE orders PARTITION BY shop, DESCRIPTOR(ts), INTERVAL '5' MINUTE)のウィンドウ・店ごとにwindow_start, window_end, shop, COUNT(*) AS cntを出す集計を書き、出力を、/root/flink/windows/session.outに保存してください。
セッションは、同じ店の隣り合う注文の間隔が5分以下ならつながり、超えると新しいセッションが始まります。セッションの開始は最初の注文の時刻、終わりは最後の注文 + 5分です。PARTITION BYを抜くと、店を混ぜてセッションを分けます。
ウィンドウ集計の上のウィンドウTop-N
/root/flink/windows/top-shops.sqlに、10分のTUMBLEウィンドウ・店別の売上(SUM(amount) AS revenue)を求めたあと、ROW_NUMBER() OVER (PARTITION BY window_start, window_end ORDER BY revenue DESC) AS rownumでウィンドウごとに上位2店だけを残して、window_start, window_end, shop, revenue, rownumを出してください。出力は、/root/flink/windows/top-shops.outに保存します。
ウィンドウ集計をサブクエリにして、その上でROW_NUMBERを付けたあと、外側でrownum <= 2で絞り込みます。PARTITION BYにウィンドウの列2つがあってはじめてウィンドウTop-Nになり、ウィンドウが閉じるときに1回だけ結果を出します。店が1つしかないウィンドウは、1行だけ出ます。
ウィンドウTVFの直上のウィンドウTop-N
/root/flink/windows/top-orders.sqlに、集計なしで、30分のTUMBLEウィンドウTVFの直上でROW_NUMBER() OVER (PARTITION BY window_start, window_end ORDER BY amount DESC)で順位を付けて、ウィンドウごとに金額が大きい注文3件のorder_id, shop, amount, window_start, window_end, rownumを出してください。出力は、/root/flink/windows/top-orders.outに保存します。
ウィンドウTop-Nは、ウィンドウ集計なしでも、ウィンドウTVFの結果の上に直接載せられます。このとき、順位の対象は集計行ではなく、注文の行そのものです。GROUP BYは使いません。金額に同点がないように作ったデータなので、順位が揺れません。
報告書: ウィンドウごとに何回数えられるか
/root/flink/windows/report.jsonに、orders(tumble.outのcntの合計)、hop_assignments(hop.outのcntの合計)、cumulate_windows(cumulate.outの行数)、sessions(session.outの行数)、max_session_orders(session.outのcntの最大値)を、整数で書いてください。
結果の行は「| +I |」で始まります。awk -F'|'で分けると、1つ目のフィールドは空で、2つ目のフィールドがopで、その後ろにSELECTの列の順に続きます。hop_assignmentsをordersと比べてみてください。size/slide倍です。