Apache Flink — ストリームを本物のエンジンで動かす
動的テーブルと変更ログ — ストリーミングの結果は書き換わり続ける表
一言でいうと
Flink SQLでは、ストリームは行が追加され続けるテーブル(動的テーブル)であり、その上のクエリの結果もテーブルです。結果のテーブルが変わるたびに、エンジンはチェンジログ(+Iは追加、-Uは古い値のリトラクション、+Uは新しい値、-Dは削除)を出力します。どのオペレーターがどの種類の変更を出すかが、そのまま、どのシンクに書けるかを決めます。
なぜ必要なのか
バッチSQLは、入力がすべてそろったあとで1回動いて終わります。結果は1度書けば済みます。ストリームには「すべてそろったあと」がありません。ユーザーごとのクリック数を数えるクエリをストリームにかけると、最初のクリックが来たときにu04 → 1と答えたあと、2回目のクリックが来たらその答えを直さなければなりません。すでに出力した答えをどう直すか。これが、ストリーミングSQLの核心の問題です。
Flinkが選んだ答えは、データベースのマテリアライズドビュー(materialized view)を借りてくることです。公式ドキュメント(Dynamic Tables)は、こうまとめています。データベースのテーブルは、INSERT・UPDATE・DELETEの流れ(changelog stream)で作られ、マテリアライズドビューは、その流れを受けて結果を書き換え続けます。ストリームをテーブルとして、継続クエリをビューの維持として見れば、SQLの意味をそのままストリームに移せます。そして、ドキュメントの約束が1つあります。継続クエリの結果は、どの瞬間でも、その時点の入力スナップショットに同じクエリをバッチで実行した結果と、意味が同じです。このモジュールのラボは、この約束を自分で確認するところから始まります。
どう動くのか
ストリーム → テーブルの場合、ソースの各レコードは、結果のテーブルに対するINSERTとして解釈されます。ファイルやログから読んだクリックは、追加だけが行われる(append-only)テーブルです。
テーブル → テーブルの場合、クエリは入力テーブルが変わるたびに、結果のテーブルを書き換えます。フィルター(WHERE dwell > 60)や射影は、入力の1行が結果の1行になって終わりなので、結果も追加のみです。一方、GROUP BY user_idのCOUNT(*)は、同じキーの行が来るたびに、すでに出力した結果の行を変えなければなりません。ドキュメントは前者をappendクエリ、後者をupdateクエリと呼び、updateクエリは、すでに出力した結果を直すためにより多くの状態を保持しなければならないと記しています。
テーブル → ストリームの場合、結果のテーブルの変化を外に出力するときは、エンコーディングを選ぶ必要があります。ドキュメントが挙げる3つです。
| エンコーディング | 何を送るか | 受け取る側に必要なもの |
|---|---|---|
| append-only | 追加された行だけ | なし。結果が追加のみのときだけ可能 |
| retract | INSERTは追加、DELETEはリトラクション、UPDATEはリトラクション(古い行)+追加(新しい行)の2つのメッセージ | なし。古い行の値を受け取って、そのまま削除する |
| upsert | INSERT・UPDATEはupsertの1つのメッセージ、DELETEは削除 | 一意キー。キーで探して上書きする |
sql-clientの表モードの出力に見えるop列が、まさにこの変更の種類です。+Iは追加、-Uは更新前の値のリトラクション、+Uは更新後の値、-Dは削除です。ユーザーごとのCOUNTなら、新しいユーザーは+Iの1行、既存のユーザーは-U・+Uの2行を出します。入力が追加のみで並列度が1なら、このログの順序と件数は、入力の順序で完全に決まります。ラボの採点ツールが、ログをPythonで1行ずつ再現して突き合わせる根拠です。
集計の上に集計を重ねると、削除も生まれます。「クリック数がn回のユーザーが何人か」を数えると、1人のユーザーがクリック数1のグループからクリック数2のグループへ移るとき、内側の集計の-Uが、外側のクリック数1のグループを1つ減らします。そのグループが0になると、結果の行そのものが消えなければならないので、-Dが出ます。
オペレーターごとにどんな変更を出すかは、EXPLAIN CHANGELOG_MODEが実行計画に表示してくれます。
GroupAggregate(groupBy=[user_id], ..., changelogMode=[I,UB,UA]) ← 갱신 전·후를 다 낸다
GroupAggregate(groupBy=[user_id], ..., changelogMode=[I,UA]) ← 같은 집계, 업서트 싱크 앞
Calc(select=[user_id, url], where=[>(dwell, 60)], changelogMode=[I])
同じ集計なのにモードが2つある理由は、シンクが何を受け取れるかを、オプティマイザーが逆向きにさかのぼって決めるからです。キーのないprintシンクは、リトラクションがあってはじめて古い行を消せるので、UBを要求し、PRIMARY KEYが宣言されたシンクは、キーで上書きすればよいので、UBを省きます。メッセージ数がほぼ半分に減ります。逆に、filesystemシンクのように、1度書いた行を直せないシンクにupdateクエリを入れると、ジョブが起動する前に「doesn't support consuming update changes」で拒否されます。それでもログをファイルに残したいなら、TO_CHANGELOGで変更の種類を列として取り出して、すべての行を追加に変えればよいです(ドキュメントのChangelog Conversion)。
現場での姿
最もよくある事故は、ダッシュボードの数字が2倍に膨らんでいることです。リトラクトストリームを受ける側が-Uを無視して+Uだけを足したか、逆に追加しかしないストレージにログをそのまま貼り付けたからです。opを見ずに行だけを数えると、更新がそのまま重複になります。結果テーブルの定義(キーが何か)と、シンクの書き込み方式(上書きか追加か)を、組にして確認するのが先です。
2つ目は、「ストリーミングの結果がバッチと違う」という報告です。最終状態を比較すると、たいてい同じです。違うのは、途中のログを誰がどう受け取ったかです。ただし、結果が同じという保証は、入力が同じときの話です。時間によって行を捨てるクエリ(次のモジュールのウォーターマーク)や、状態を時間で消す設定が加わると、バッチと変わることがあります。
3つ目は、アップサートシンクのキーが、結果テーブルのキーとずれている場合です。集計キーはuser_idなのに、シンクのPRIMARY KEYを別の列にすると、シンクが受け取る変更が、シンクのキー基準で順序が入れ替わることがあります。ドキュメント(Configurationのtable.exec.sink.upsert-materialize)は、このような入れ替わりが起きると、オプティマイザーがシンクの前にupsert materializeオペレーターを挟むと記しています。キーごとに受け取った行を状態として保持しておき、正しいupsertの順序で出し直すオペレーターです。状態がもう1つ増えることになるので、シンクのキーは最初から結果テーブルの一意キーと同じにするのが原則です。
次のラボですること
クリック120件をユーザーごとに数えるクエリを、バッチとストリーミングでそれぞれ動かして、最終状態が同じであることと、ストリーミングのログが229行であることを確認します。集計の上の集計で-Dが出ることを見て、EXPLAIN CHANGELOG_MODEで集計とフィルターの変更の種類を比べます。集計をファイルシンクに入れて拒否されたあと、TO_CHANGELOGに変えてファイルに書き、キーのないシンクとキーのあるシンクの前で、計画がどう変わるかを確認します。最後に、同じログをアップサートで送っていたなら、メッセージが何件だったかを報告書に書きます。