Apache Flink — ストリームを本物のエンジンで動かす
ストリーム結合 — 何をどれだけ長く覚えておくか
一言でいうと
ストリームの結合は、どちら側の行をどれだけ長く記憶するかで分かれます。時間条件のない通常結合は、両側を永遠に保持し、区間結合は、時間の範囲が過ぎると捨て、イベント時間のテンポラル結合(temporal join)は、右側のバージョンだけを保持して、左側の行を「その時刻のバージョン」1つと結びます。
なぜ必要なのか
バッチでは、結合は簡単です。2つのテーブルがすでにすべてそろっているので、片方でハッシュテーブルを作り、もう片方を走査すれば終わりです。ストリームでは、2つのテーブルが終わりません。注文が入った瞬間には、その注文の配送はまだなく、ユーザー情報は、ずっと前に届いていることも、ずっとあとに届くこともあります。そのため、エンジンは「今届いた行と対応しうる反対側の行」をどこかに積んでおく必要があり、その山がそのまま状態です。
問題は、いつ捨てるかです。対応する行が永遠に来ないかもしれない行を無期限に保持すると、状態が際限なく増えます。かといって、いつでも捨ててしまうと、遅れて届いた対応する行を逃します。Flink SQLが結合を複数の種類に分けている理由がこれです。クエリが時間について何を約束するかによって、エンジンが安全に捨てられるものが変わります。
どう動くのか
通常結合(regular join)は、最も自由です。公式ドキュメントの表現どおり、片方に新しい行が来ると、反対側の過去と未来のすべての行と照らし合わせます。そのため、午後に登録したユーザーの午前の注文も結合されます。時間を見ないからです。代償も、ドキュメントにそのまま書かれています。両側の入力を永遠に状態に置く必要があります。状態TTLで減らすことはできますが、そうすると結果が間違うことがあります。実行計画をJSONで取り出してみると(COMPILE PLAN)、結合ノードにleftState・rightStateの2つの状態が、TTL 0 ms(削除しない)とともに書かれています。
区間結合(interval join)は、等号条件1つと、両側の時間を結ぶ範囲を要求します。s.ship_time BETWEEN o.order_time AND o.order_time + INTERVAL '2' HOURがその例です。入力は、時間属性のあるappend-onlyのテーブルでなければなりません。時間属性はほぼ単調に増えるので、ウォーターマークが(注文時刻 + 2時間)を過ぎると、その注文と対応する配送はもう来ないと確定して、状態から消します。実測で確認した境界もあります。BETWEENは両端を含むので、遅延0秒とちょうど7200秒の配送は結合され、7201秒は結合されません。LEFT JOINに変えて時間条件をONに置くと、範囲が閉じるまで対応する行が見つからなかった注文が、nullとともに1回出ます。結果は最後までappend-onlyです。
イベント時間のテンポラル結合(temporal join)は、左側(注文)の行1つを、右側のバージョン付きテーブルの「その時刻に有効だったバージョン」1つと結びます。構文は、SQL:2011のFOR SYSTEM_TIME AS OF o.order_timeです。バージョン付きテーブルになるには、主キーとイベント時間属性が必要です。為替レートのファイルのようにappend-onlyのソースには、主キーを設定できませんが、ドキュメントはここでコツを教えています。通貨ごとにROW_NUMBER() OVER (PARTITION BY currency ORDER BY update_time DESC) = 1の重複排除ビューを作ると、オプティマイザーがcurrencyを主キーと推論して、バージョン付きビューとして使います。
CREATE TEMPORARY VIEW rates_v AS
SELECT currency, rate, update_time FROM (
SELECT *, ROW_NUMBER() OVER (PARTITION BY currency ORDER BY update_time DESC) AS rn FROM rates)
WHERE rn = 1;
SELECT o.order_id, r.rate
FROM orders o JOIN rates_v FOR SYSTEM_TIME AS OF o.order_time AS r
ON o.currency = r.currency;
実測の結果、注文は、update_time <= order_timeのバージョンのうち最も遅いものと結合されました(為替レートが注文と同じ秒に変わるなら、新しい為替レート)。最初の為替レートより早い注文は、結合するバージョンがなく、内部結合から外れました。ドキュメントのとおり、この結合は両側のウォーターマークがトリガーで、右側があとから変わっても、すでに出した結果を修正しません。古いバージョンは、必要なくなると状態から消えます。計画JSONを見ると、テンポラル結合のノードと区間結合のノードには、TTLが付いたstateの項目がそもそもありません。この2つは、TTLではなく、時間で状態を整理します。
同じバージョン付きビューをFOR SYSTEM_TIME AS OFなしで結合すると、ただの通常結合です。為替レートが変わるたびに、過去の注文まで新しい為替レートで再計算されて、リトラクション(-U)と更新(+U)が大量に出て、最終結果は、最後の為替レートで換算した値になります。
| 結合 | 記憶するもの | 結果を修正するか | 状態を消す根拠 |
|---|---|---|---|
| 通常 | 両側のすべて | 修正する(更新入力なら) | TTLのみ(正確さを失うことがある) |
| 区間 | 時間範囲の中の行 | 修正しない | ウォーターマークが範囲の終わりを過ぎた |
| イベント時間のテンポラル | 右側の必要なバージョン | 修正しない | ウォーターマークが過ぎて不要になったバージョン |
現場での姿
最もよくある事故は、「注文に商品情報を結合しただけなのに、状態が何か月も増え続ける」です。原因は、ほとんどいつも通常結合です。商品テーブルは小さくても、注文側が永遠に積み上がります。ディメンション情報が、その時点の値であればよいなら、テンポラル結合が正しい道具です。
2つ目は、売上の再計算事故です。為替レートや価格表を通常結合で結合しておくと、価格が変わった瞬間に、過去の注文の金額が静かに変わります。ダッシュボードの昨日の売上が今日は変わっているという報告が来たら、結合の種類から確認します。このラボで、同じUSDの注文を2つの方式で換算すると、合計が実際に変わります。
3つ目は、区間結合の境界です。「2時間以内に配送」を<で書くかBETWEENで書くかによって、ちょうど2時間で出た配送が分かれます。運用指標の定義書とSQLの不等号を、1度は突き合わせる必要があります。
次のラボですること
4つのソースを定義し、通常結合でティアごとの注文を集めます。区間結合で2時間以内の配送を結合し、外部区間結合で期限内に配送されなかった注文を探します。為替レートをバージョン付きビューにして、テンポラル結合で注文時刻の為替レートを結合し、同じビューを通常結合で結合したとき、合計がどう変わるかを見ます。3つの結合の実行計画をJSONで取り出して、状態の項目を比べ、数字を報告書にまとめます。