TT Lab
はじめる
学ぶ 学習パス コース

Apache Flink — ストリームを本物のエンジンで動かす

4種類のウィンドウ TVF で同じ注文を切り分ける

TT Labで続きを見る

目標

同じ注文ストリームにTUMBLE・HOP・CUMULATE・SESSIONウィンドウをかけて、ウィンドウの列がどう付き、1つの行が何個のウィンドウに入るかを、結果で確認します。ウィンドウ集計の上のウィンドウTop-Nと、ウィンドウTVFの直上のウィンドウTop-Nを作ります。

なぜ重要なのか

ウィンドウの種類を間違えて選ぶと、数字が静かに間違います。HOPの件数を合計すると重なった分だけ膨らみ、GROUP BYからウィンドウの列を抜くと、ウィンドウ集計ではなく無限の集計になって、更新ログが混ざります。ウィンドウTVFは、行にウィンドウの列を付ける関数にすぎないので、その列がどう付くかさえ正確に知っていれば、集計でも順位でも、普通のSQLで載せられます。このラボの採点ツールは、クラスターに問い合わせません。皆さんが保存したsql-clientの出力を読み、ウィンドウごとの期待値を、元のCSVからPythonで直接計算して照合します。

ステップ

  1. 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に保存してください。
  2. /root/flink/windows/tumble.sqlに、10分のTUMBLEウィンドウ・店(shop)ごとにcnt(件数)・revenue(amountの合計)を出すクエリを書いて実行し、出力を、/root/flink/windows/tumble.outに保存してください。
  3. /root/flink/windows/hop.sqlに、HOP(slide 5分、size 10分)のウィンドウごとにcntを出すクエリを書いて実行し、出力を、/root/flink/windows/hop.outに保存してください。
  4. /root/flink/windows/cumulate.sqlに、CUMULATE(step 10分、size 1時間)のウィンドウごとにrevenueを出すクエリを書いて実行し、出力を、/root/flink/windows/cumulate.outに保存してください。
  5. /root/flink/windows/session.sqlに、SESSION(PARTITION BY shop、gap 5分)のウィンドウ・店ごとにcntを出すクエリを書いて実行し、出力を、/root/flink/windows/session.outに保存してください。
  6. /root/flink/windows/top-shops.sqlに、10分のTUMBLEウィンドウごとに、店別の売上の上位2店(window_start, window_end, shop, revenue, rownum)を出すクエリを書いて実行し、出力を、/root/flink/windows/top-shops.outに保存してください。
  7. /root/flink/windows/top-orders.sqlに、30分のTUMBLEウィンドウごとに、金額が大きい注文3件(order_id, shop, amount, window_start, window_end, rownum)を集計なしで出すクエリを書いて実行し、出力を、/root/flink/windows/top-orders.outに保存してください。
  8. /root/flink/windows/report.jsonに、orders・hop_assignments・cumulate_windows・sessions・max_session_ordersを書いてください。

参考

ウィンドウ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倍です。