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

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

ストリーム結合 — 何をどれだけ長く覚えておくか

TT Labで続きを見る

一言でいうと

ストリームの結合は、どちら側の行をどれだけ長く記憶するかで分かれます。時間条件のない通常結合は、両側を永遠に保持し、区間結合は、時間の範囲が過ぎると捨て、イベント時間のテンポラル結合(temporal join)は、右側のバージョンだけを保持して、左側の行を「その時刻のバージョン」1つと結びます。

なぜ必要なのか

バッチでは、結合は簡単です。2つのテーブルがすでにすべてそろっているので、片方でハッシュテーブルを作り、もう片方を走査すれば終わりです。ストリームでは、2つのテーブルが終わりません。注文が入った瞬間には、その注文の配送はまだなく、ユーザー情報は、ずっと前に届いていることも、ずっとあとに届くこともあります。そのため、エンジンは「今届いた行と対応しうる反対側の行」をどこかに積んでおく必要があり、その山がそのまま状態です。

問題は、いつ捨てるかです。対応する行が永遠に来ないかもしれない行を無期限に保持すると、状態が際限なく増えます。かといって、いつでも捨ててしまうと、遅れて届いた対応する行を逃します。Flink SQLが結合を複数の種類に分けている理由がこれです。クエリが時間について何を約束するかによって、エンジンが安全に捨てられるものが変わります。

どう動くのか

3つのコマの図。1つ目のコマの通常結合は、注文とユーザーの両側の行がすべて状態に積まれ、どちら側に新しい行が来ても、反対側の全体と照らし合わせます。2つ目のコマの区間結合は、注文時刻から2時間幅の帯の中にある配送だけが結合され、ウォーターマークが帯の終わりを過ぎると、その注文は状態から消えます。3つ目のコマのテンポラル結合は、為替レートが階段状に変わるバージョンの列で、注文は自分の時刻に有効だった階段の1段とだけ結合されます。最初のバージョンより早い注文は、結合するバージョンがありません

通常結合(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で取り出して、状態の項目を比べ、数字を報告書にまとめます。