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

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

重複を取り除き、順位を付ける

TT Labで続きを見る

目標

同じROW_NUMBERのパターンで、先頭行を残す方式・末尾行を残す方式・Top-Nを作って、それぞれが出すチェンジログの形と量を、出力件数と実行計画で確認します。

なぜ重要なのか

再送された出来事を取り除くときと、現在の状態を残すときは、SQLではASCとDESCの1語の違いですが、前者はappend-onlyの結果を、後者はリトラクションと更新のログを出します。この違いが、どのシンクを使えるか、下流の集計が何を受け取るかを決めます。Top-Nは、順位番号を出力するかによって、同じ結果でもログの量が大きく変わります。このラボの採点ツールは、クラスターに問い合わせず、皆さんが保存したsql-clientの出力を、元のCSVから直接計算した値と照合します。

ステップ

  1. flink-upでクラスターを起動し、/root/flink/dedup/ddl.sqlに、events・salesを定義してください(eventsのtsにウォーターマーク)。/root/flink/dedup/count.sqlに、バッチでn_rows(全体の行数)・n_orders(異なる注文の数)を出すクエリを書いて実行し、出力を、/root/flink/dedup/count.outに保存してください。
  2. /root/flink/dedup/first.sqlに、注文ごとに最も早いイベントを残す、先頭行を残す方式のクエリ(order_id・status・amount)を書いて実行し、出力を、/root/flink/dedup/first.outに保存してください。
  3. /root/flink/dedup/last.sqlに、注文ごとに最も遅いイベントを残す、末尾行を残す方式のクエリを書いて実行し、出力を、/root/flink/dedup/last.outに保存してください。
  4. /root/flink/dedup/explain.sqlに、ステップ2・3のクエリをEXPLAIN CHANGELOG_MODEするクエリを書いて実行し、出力を、/root/flink/dedup/explain.outに保存してください。
  5. /root/flink/dedup/board.sqlに、末尾行を残す方式の上でstatusごとにorders(注文数)・amount(合計)を数えるクエリを書いて実行し、出力を、/root/flink/dedup/board.outに保存してください。
  6. /root/flink/dedup/topn.sqlに、salesでカテゴリごとの累積qtyの合計の上位3つを、順位番号とともに出すクエリ(category・product・total・rn)を書いて実行し、出力を、/root/flink/dedup/topn.outに保存してください。
  7. /root/flink/dedup/topn-norank.sqlに、同じTop-3で、外側のSELECTからrnだけを除いたクエリを書いて実行し、出力を、/root/flink/dedup/topn-norank.outに保存してください。
  8. /root/flink/dedup/report.jsonに、n_rows・n_orders・update_pairs・shipped_now・topn_log_rows・norank_log_rows・top_booksを書いてください。

参考

2つのソースを定義して行数を数える

flink-upのあと、/root/flink/dedup/ddl.sqlに、events・salesを定義してください(eventsのtsにWATERMARK FOR ts AS ts)。/root/flink/dedup/count.sqlに、バッチモードでn_rows・n_ordersを出すクエリを書き、sql-client.sh -i ddl.sql -f count.sqlで実行して、出力を、/root/flink/dedup/count.outに保存してください。

注文数はCOUNT(DISTINCT order_id)です。2つの値の差が、そのまま、同じ注文について2番目以降に来たイベントの数です。ステップ3で、この数字にもう一度出会います。

先頭行を残す: 1回出したら終わり

/root/flink/dedup/first.sqlに、ストリーミングモードで注文(order_id)ごとにtsが最も早いイベントを1つ残す重複排除を書いて、order_id・status・amountを出してください。出力は、/root/flink/dedup/first.outに保存します。

内側のSELECTでROW_NUMBER() OVER (PARTITION BY order_id ORDER BY ts ASC)で順位番号を付け、外側で順位番号が1のものだけを選びます。出力のopがすべて+Iになっているか見てください。1度出した先頭行は、変わることがありません。

末尾行を残す: リトラクションして直す

/root/flink/dedup/last.sqlに、注文ごとにtsが最も遅いイベントを残す重複排除を書き(order_id・status・amount)、出力を、/root/flink/dedup/last.outに保存してください。ログの件数(+I・-U・+U)が、ステップ1の2つの数字とどんな関係にあるかを見てください。

並べ替えの方向を変えるだけです。同じ注文のイベントが新しく来るたびに、先に出した行を-Uで取り下げて、新しい行を+Uで出します。再送でまったく同じ行が来ても、組が出ます。処理時間で並べると、この件数が揺れることがあるので、tsを使ってください。

計画で2つの重複排除を分ける

/root/flink/dedup/explain.sqlに、ステップ2とステップ3のクエリの前に、それぞれEXPLAIN CHANGELOG_MODEを付けた2つの文を書き、出力を、/root/flink/dedup/explain.outに保存してください。

実行計画(Optimized Execution Plan)にDeduplicate(keep=[...])ノードが、物理計画(Optimized Physical Plan)の行ごとにchangelogMode=[...]が付きます。2つのクエリで、keepとoutputInsertOnly、changelogModeがどう違うかを比べてください。

現在の状態で数える: 重複排除の上の集計

/root/flink/dedup/board.sqlに、ステップ3の末尾行を残す方式を内側に置き、外側でGROUP BY statusでorders(注文数)・amount(amountの合計)を出すクエリを書き、出力を、/root/flink/dedup/board.outに保存してください。

注文1つがcreatedからpaidに変わると、重複排除が-U created・+U paidを出し、集計は、createdのグループから引いてpaidのグループに足します。最終的なordersの合計は、注文数と同じになるはずです。重複排除なしでeventsを直接数えると、この合計が行数の分だけ大きくなります。

Top-3: 順位が変わるたびに直す

/root/flink/dedup/topn.sqlに、salesを(category, product)ごとにSUM(qty) AS totalでまとめたあと、カテゴリごとにtotalの降順で上位3つを、順位番号rnとともに出すクエリを書いてください(category・product・total・rn)。出力は、/root/flink/dedup/topn.outに保存します。

GROUP BYの合計をサブクエリにして、その上にROW_NUMBER() OVER (PARTITION BY category ORDER BY total DESC)を付け、外側でrn <= 3で選びます。合計が変わるたびに順位が動くので、出力に-U・-Dが混ざって出ます。最終状態だけを見れば、カテゴリごとに3行です。

順位番号を除くとログが減る

/root/flink/dedup/topn-norank.sqlに、ステップ6のクエリから外側のSELECTのrnだけを除いたクエリ(category・product・total)を書いて実行し、出力を、/root/flink/dedup/topn-norank.outに保存してください。最終的な上位3つは同じで、ログの行数はステップ6より少なくなければなりません。

順位番号の列は結果の一意キーの一部なので、出力すると、1つの商品の順位が上がるとき、その下の順位の行がすべて送り直されます。順位番号を除けば、変わった商品の行だけを送れば済みます。2つのファイルのログの行数を、grep -cで数えて比べてみてください。

報告書: ログの量を数字で

/root/flink/dedup/report.jsonに、n_rows・n_orders(ステップ1)、update_pairs(last.outの-Uの数)、shipped_now(board.outの最終状態のshipped注文数)、topn_log_rows・norank_log_rows(2つのTop-3の出力のログ行数)、top_books(booksカテゴリの1位の商品)を書いてください。

すべて、前に保存した出力から写します。ログの行は、「| +I |」のようにopで始まる行です。チェンジログがあるテーブルの最終値は、そのキーの最後の+I・+Uです。update_pairsは、ステップ1の2つの数字でも合っているかを確認してみてください。