Apache Flink — ストリームを本物のエンジンで動かす
重複を取り除き、順位を付ける
目標
同じROW_NUMBERのパターンで、先頭行を残す方式・末尾行を残す方式・Top-Nを作って、それぞれが出すチェンジログの形と量を、出力件数と実行計画で確認します。
なぜ重要なのか
再送された出来事を取り除くときと、現在の状態を残すときは、SQLではASCとDESCの1語の違いですが、前者はappend-onlyの結果を、後者はリトラクションと更新のログを出します。この違いが、どのシンクを使えるか、下流の集計が何を受け取るかを決めます。Top-Nは、順位番号を出力するかによって、同じ結果でもログの量が大きく変わります。このラボの採点ツールは、クラスターに問い合わせず、皆さんが保存したsql-clientの出力を、元のCSVから直接計算した値と照合します。
ステップ
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に保存してください。- /root/flink/dedup/first.sqlに、注文ごとに最も早いイベントを残す、先頭行を残す方式のクエリ(
order_id・status・amount)を書いて実行し、出力を、/root/flink/dedup/first.outに保存してください。 - /root/flink/dedup/last.sqlに、注文ごとに最も遅いイベントを残す、末尾行を残す方式のクエリを書いて実行し、出力を、/root/flink/dedup/last.outに保存してください。
- /root/flink/dedup/explain.sqlに、ステップ2・3のクエリを
EXPLAIN CHANGELOG_MODEするクエリを書いて実行し、出力を、/root/flink/dedup/explain.outに保存してください。 - /root/flink/dedup/board.sqlに、末尾行を残す方式の上で
statusごとにorders(注文数)・amount(合計)を数えるクエリを書いて実行し、出力を、/root/flink/dedup/board.outに保存してください。 - /root/flink/dedup/topn.sqlに、salesでカテゴリごとの累積
qtyの合計の上位3つを、順位番号とともに出すクエリ(category・product・total・rn)を書いて実行し、出力を、/root/flink/dedup/topn.outに保存してください。 - /root/flink/dedup/topn-norank.sqlに、同じTop-3で、外側のSELECTから
rnだけを除いたクエリを書いて実行し、出力を、/root/flink/dedup/topn-norank.outに保存してください。 - /root/flink/dedup/report.jsonに、
n_rows・n_orders・update_pairs・shipped_now・topn_log_rows・norank_log_rows・top_booksを書いてください。
参考
- 元データ(ヘッダーなしのCSV、時刻は秒単位、行の順序が到着順で、tsは昇順):
dedup_events.csv=order_id, status, amount, ts。1つの注文に1–4件あり、一部は同じ行がすぐ後ろにもう1回来ます(再送)。dedup_sales.csv=sale_id, category, product, qty, ts。すべて/opt/lab/fixtures/data/にあります。 - 実行:
sql-client.sh -i ddl.sql -f 쿼리.sql > 쿼리.out 2>&1(プレースホルダーはクエリ名です)。ストリーミングの結果は、先頭にop列(+I・-U・+U・-D)が付きます。 - 重複排除は、ドキュメントの形を正確に守る必要があります。内側のSELECTに
ROW_NUMBER() OVER (PARTITION BY 키 ORDER BY 시간속성 ASC|DESC) AS rn(プレースホルダーは順に、キーと時間属性です)を、外側にWHERE rn = 1を置きます。 - よくある間違い: 処理時間(
PROCTIME())で並べると、結果が到着時刻によって変わることがあります。このラボは、イベント時間tsで並べます。 - 公式ドキュメント: Deduplication・Top-N・EXPLAIN・Dynamic Tables
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つの数字でも合っているかを確認してみてください。