Apache Flink — ストリームを本物のエンジンで動かす
3種類の結合で注文をつなぐ
目標
同じ注文の流れに、ユーザー・配送・為替レートを、通常結合・区間結合・イベント時間のテンポラル結合で結合してみて、それぞれの結合が何を記憶して結果を修正するかを、出力と実行計画で確認します。
なぜ重要なのか
ストリームの結合は、対応する行がいつ来るかわからないので、反対側の行を状態に積んでおきます。時間条件がなければ永遠に、時間範囲を与えればウォーターマークが過ぎるまで、バージョン付きテーブルなら必要なバージョンだけを保持します。結合の種類を間違えて選ぶと、状態が際限なく増えたり、過去の結果が静かに変わったりします。このラボの採点ツールは、クラスターに問い合わせず、皆さんが保存したsql-clientの出力と計画JSONを、元のCSVから直接計算した値と照合します。
ステップ
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に保存してください。- /root/flink/joins/regular.sqlに、ordersとusersを
user_idで結合して、ティアごとにorders(件数)・amount(合計)を出すストリーミングクエリを書き、出力を、/root/flink/joins/regular.outに保存してください。 - /root/flink/joins/interval.sqlに、注文時刻から2時間以内(両端を含む)に出た配送を結合する区間結合を書き、
order_id・ship_id・delay_s(秒)を、/root/flink/joins/interval.outに保存してください。 - /root/flink/joins/unshipped.sqlに、2時間以内に配送されなかった注文の
order_id・order_timeを出す外部区間結合を書き、出力を、/root/flink/joins/unshipped.outに保存してください。 - /root/flink/joins/temporal.sqlに、ratesでバージョン付きビュー
rates_vを作成し、注文時刻の為替レートを結合するテンポラル結合で、order_id・currency・rate・amount_krw(amount × rate)を出すクエリを書いて実行し、出力を、/root/flink/joins/temporal.outに保存してください。 - /root/flink/joins/latest.sqlに、同じ
rates_vをFOR SYSTEM_TIME AS OFなしで結合して、通貨ごとにorders・total_krwを出すクエリを書いて実行し、出力を、/root/flink/joins/latest.outに保存してください。 - 3つの結合(ステップ2・3・5)の実行計画を
COMPILE PLANで、/root/flink/joins/regular-plan.json・/root/flink/joins/interval-plan.json・/root/flink/joins/temporal-plan.jsonに取り出してください。 - /root/flink/joins/report.jsonに、
matched_orders・interval_rows・unshipped_orders・temporal_rows・orders_without_rate・usd_gap_krwを書いてください。
参考
- 元データ(ヘッダーなしのCSV、時刻は秒単位、ファイルごとに自分の時間順に並んでいます):
joins_users.csv=user_id, tier, signup_time・joins_orders.csv=order_id, user_id, currency, amount, order_time・joins_shipments.csv=ship_id, order_id, ship_time・joins_rates.csv=currency, rate, update_time(rateはウォン単位で、DECIMAL(10, 4)で読んでください)。すべて/opt/lab/fixtures/data/にあります。 - 定義を1回だけ書くには、
sql-client.sh -i ddl.sql -f 쿼리.sql > 쿼리.out 2>&1(プレースホルダーはクエリ名です)を使います。-iのファイルのCREATE文が、先に実行されます。 - ストリーミングの結果は、先頭に
op列(+I・-U・+U・-D)が付きます。ジョブはファイルの終わりで終了し、そのときウォーターマークが最後まで進んで、残った範囲・バージョンがすべて処理されます。 COMPILE PLANは、INSERT INTO文を受け取ります。'connector' = 'blackhole'のシンクを1つ作って使ってください。同じパスにファイルがあると、上書きせずにエラーになるので、取り出し直すときは先に削除します。- よくある間違い: 外部区間結合の時間条件を
WHEREに置くと、null行が除外されます。テンポラル結合の右側には主キーが必要ですが、append-onlyのソースは、重複排除ビューにする必要があります。 - 公式ドキュメント: Joins・Versioned Tables・Deduplication・SQL Client
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は数値で書きます。