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

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

3種類の結合で注文をつなぐ

TT Labで続きを見る

目標

同じ注文の流れに、ユーザー・配送・為替レートを、通常結合・区間結合・イベント時間のテンポラル結合で結合してみて、それぞれの結合が何を記憶して結果を修正するかを、出力と実行計画で確認します。

なぜ重要なのか

ストリームの結合は、対応する行がいつ来るかわからないので、反対側の行を状態に積んでおきます。時間条件がなければ永遠に、時間範囲を与えればウォーターマークが過ぎるまで、バージョン付きテーブルなら必要なバージョンだけを保持します。結合の種類を間違えて選ぶと、状態が際限なく増えたり、過去の結果が静かに変わったりします。このラボの採点ツールは、クラスターに問い合わせず、皆さんが保存したsql-clientの出力と計画JSONを、元のCSVから直接計算した値と照合します。

ステップ

  1. flink-upでクラスターを起動し、/root/flink/joins/ddl.sqlに、4つのソース(users・orders・shipments・rates)を定義してください。orders・shipments・ratesの時間列には、ウォーターマークを設定します。/root/flink/joins/count.sqlに、バッチで4つのテーブルの行数を数えるクエリ(列名はtbl・n)を書いて実行し、出力を、/root/flink/joins/count.outに保存してください。
  2. /root/flink/joins/regular.sqlに、ordersとusersをuser_idで結合して、ティアごとにorders(件数)・amount(合計)を出すストリーミングクエリを書き、出力を、/root/flink/joins/regular.outに保存してください。
  3. /root/flink/joins/interval.sqlに、注文時刻から2時間以内(両端を含む)に出た配送を結合する区間結合を書き、order_id・ship_id・delay_s(秒)を、/root/flink/joins/interval.outに保存してください。
  4. /root/flink/joins/unshipped.sqlに、2時間以内に配送されなかった注文のorder_id・order_timeを出す外部区間結合を書き、出力を、/root/flink/joins/unshipped.outに保存してください。
  5. /root/flink/joins/temporal.sqlに、ratesでバージョン付きビューrates_vを作成し、注文時刻の為替レートを結合するテンポラル結合で、order_id・currency・rate・amount_krw(amount × rate)を出すクエリを書いて実行し、出力を、/root/flink/joins/temporal.outに保存してください。
  6. /root/flink/joins/latest.sqlに、同じrates_vをFOR SYSTEM_TIME AS OFなしで結合して、通貨ごとにorders・total_krwを出すクエリを書いて実行し、出力を、/root/flink/joins/latest.outに保存してください。
  7. 3つの結合(ステップ2・3・5)の実行計画をCOMPILE PLANで、/root/flink/joins/regular-plan.json・/root/flink/joins/interval-plan.json・/root/flink/joins/temporal-plan.jsonに取り出してください。
  8. /root/flink/joins/report.jsonに、matched_orders・interval_rows・unshipped_orders・temporal_rows・orders_without_rate・usd_gap_krwを書いてください。

参考

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

flink-upのあと、/root/flink/joins/ddl.sqlに、users・orders・shipments・ratesを定義してください(order_time・ship_time・update_timeにウォーターマーク)。/root/flink/joins/count.sqlに、バッチモードで4つのテーブルの行数をtbl・nの列で出すクエリを書き、sql-client.sh -i ddl.sql -f count.sqlで実行して、出力を、/root/flink/joins/count.outに保存してください。

ウォーターマークは、WATERMARK FOR time_col AS time_colのように書きます。ファイルごとに自分の時間順に並んでいるので、遅延を与えなくても、遅延行はありません。4つのSELECTをUNION ALLでつなぐと、1つのテーブルとして出ます。

通常結合: 時間を見ず、両側を記憶する

/root/flink/joins/regular.sqlに、ストリーミングモードでordersとusersをuser_idで内部結合して、tierごとにorders(件数)・amount(amountの合計)を出すクエリを書き、出力を、/root/flink/joins/regular.outに保存してください。

通常結合は、時間条件なしで等号だけを使います。u41–u44はusersにいないユーザーなので、内部結合から外れます。午後に登録したユーザー(signup_timeが注文より遅い)の注文も結合されるかを、結果で確認してみてください。通常結合は、到着順と関係なく、過去・未来のすべての対応する行を探します。

区間結合: 2時間以内の配送だけを結合する

/root/flink/joins/interval.sqlに、ordersとshipmentsをorder_idで結合し、ship_timeが注文時刻から2時間以内(両端を含む)のものだけを残す区間結合を書いて、order_id・ship_id・delay_s(遅延の秒、TIMESTAMPDIFF(SECOND, ...))を、/root/flink/joins/interval.outに保存してください。

区間結合には、等号1つと、両側の時間を結ぶ範囲が必要です。BETWEEN a AND bは両端を含みます。遅延がちょうど0秒・7200秒の配送があり、7201秒の配送もあります。2回に分けて配送された注文は、2行になります。

外部区間結合: 期限内に対応する行がなかった注文

/root/flink/joins/unshipped.sqlに、orders LEFT JOIN shipmentsで、2時間以内に配送が1つもなかった注文のorder_id・order_timeを出すクエリを書き、出力を、/root/flink/joins/unshipped.outに保存してください。

時間条件はON句に置き、対応する行がなくnullで埋められた行だけを、WHEREで選びます。配送記録がまったくない注文だけでなく、2時間を超えて出た注文も入る必要があります。null行は、ウォーターマークが(注文時刻 + 2時間)を過ぎたあとに、1回出ます。

テンポラル結合: 注文時刻の為替レート

/root/flink/joins/temporal.sqlで、ratesから通貨ごとの最新の行だけを残すバージョン付きビューrates_vを作成し、ordersをFOR SYSTEM_TIME AS OF o.order_timeで結合して、order_id・currency・rate・amount_krw(amount × rate)を出してください。出力は、/root/flink/joins/temporal.outに保存します。

ratesはappend-onlyなので、主キーを設定できません。ROW_NUMBER() OVER (PARTITION BY currency ORDER BY update_time DESC)が1の行だけを残すビューを作ると、currencyが主キー、update_timeがイベント時間のバージョン付きビューになります。内部結合なら、最初の為替レートより早い注文は除外されます。

同じビューを通常結合で: 過去の注文が再計算される

/root/flink/joins/latest.sqlで、ステップ5のrates_vを、FOR SYSTEM_TIME AS OFを付けずにordersとcurrencyで結合して、通貨ごとにorders(件数)・total_krw(amount × rateの合計)を出し、出力を、/root/flink/joins/latest.outに保存してください。

バージョン付きビューは、為替レートが変わるたびに更新を出します。通常結合は、その更新を受けて過去の注文まで再計算するので、出力に-Uが見えます。最終的な合計は、テンポラル結合の合計と違います。どの為替レートで計算されたものかを考えてみてください。

3つの結合の実行計画から状態を読む

ステップ2・3・5の結合を、blackholeシンクに入れるINSERTに変えて、COMPILE PLANで、/root/flink/joins/regular-plan.json・/root/flink/joins/interval-plan.json・/root/flink/joins/temporal-plan.jsonを作成してください(通常結合はorders⋈users、区間結合はorders⋈shipments、テンポラル結合はorders⋈rates_v)。

COMPILE PLAN 'file:///path.json' FOR INSERT INTO sink_table SELECT ...の形です。ファイルがすでにあるとエラーになるので、取り出し直すときは削除して実行します。作成後、jqでnodesのtypeとstateを見てみてください。どの結合ノードにleftState・rightStateがあり、どのノードにはないかが要点です。

報告書: 結合ごとに何が結合され、何が外れたか

/root/flink/joins/report.jsonに、matched_orders(通常結合でユーザーと結合された注文数)、interval_rows(区間結合の結果の行数)、unshipped_orders(ステップ4の行数)、temporal_rows(テンポラル結合の結果の行数)、orders_without_rate(全体の注文数 − temporal_rows)、usd_gap_krw(latest.outのUSDのtotal_krw − temporal.outのUSDのamount_krwの合計、小数)を書いてください。

すべて、前に保存した出力から写せます。チェンジログがあるテーブルは、最後の+I/+Uが最終値です。grepで「| +I |」のような行だけを選び、awk -F'|'でフィールドを切り出せば、数えられます。整数の項目は整数で、usd_gap_krwは数値で書きます。