Apache Flink — ストリームを本物のエンジンで動かす
状態はどこにどれだけ残るか — 演算子・TTL・バックエンド
一言でいうと
状態は、ジョブ全体ではなくオペレーターごとに別々にあります。どのオペレーターが何を保持していて、TTLがいくつなのかは、COMPILE PLANがJSONで正確に書き出してくれます。際限のないGROUP BYと通常結合は、デフォルトで永遠に保持し、ウィンドウ集計は、ウィンドウが閉じるときに自分で空にします。ステートバックエンドは、その状態をどこにどんな形で置くかを決めるだけで、結果は変えません。
なぜ必要なのか
ストリーミングジョブを数週間動かすと、ほぼ必ず受ける質問が、「メモリ(またはディスク)がなぜ増え続けるのか」です。公式ドキュメントは、GROUP BYについて次のように警告しています。ストリーミングクエリの結果を計算するために必要な状態は、際限なく増える可能性があり、サイズはグループの数と集計関数の種類で決まります(MIN・MAXは重く、COUNTは軽い)。前のモジュールの通常結合は、両側の入力を永遠に保持すると述べました。
そこでTTL(アイドル状態の保持時間)を設定します。ところが、TTLはただではありません。ドキュメントの決定性の章(Determinism)は、TTLで状態を消すことが、しばしば必要な妥協である一方で、結果を非決定的にしうると記しています。消されたキーに遅れて行が来ると、初めて見るキーのように、もう一度数え始めます。ですから、「どのオペレーターにいくつを設定するか」を、オペレーター単位で把握して決める必要があります。このモジュールは、その地図を実行計画から読む方法です。
どう動くのか
COMPILE PLAN 'file:///경로.json' FOR INSERT INTO ...(プレースホルダーはパスです)は、ジョブを動かさずに、最適化された計画をJSONで書き出します。ノードごとにtypeがあり、状態のあるノードにはstateのリストが付きます。このラボ環境で取り出した結果は、次のとおりです。
| ノード(type) | stateの項目 | デフォルトのTTL |
|---|---|---|
| stream-exec-calc(フィルターと変換) | なし | — |
| stream-exec-group-aggregate(際限のないGROUP BY) | groupAggregateState | 0 ms |
| stream-exec-join(通常結合) | leftState・rightState | 0 ms |
| stream-exec-deduplicate・stream-exec-rank | deduplicateState・rankState | 0 ms |
| 区間結合・テンポラル結合・ウィンドウ集計 | なし | — |
0 msは「削除しない」という意味です。ドキュメントの設定の説明(table.exec.state.ttl)のとおりです。更新されていない状態を最低でもこれだけの間は保持し、それより長く休んだら、そのあといつか削除します。デフォルト値の0は、永遠に置くという意味です。SET 'table.exec.state.ttl' = '30 min'を指定すると、上の表の0 msがすべて30 minに変わります。表の最後の行のオペレーターは、TTLではなく時間(ウォーターマークがウィンドウの終わり・区間の終わりを過ぎること)で状態を整理するので、計画にTTLの項目がありません。ドキュメントも、「短命なウィンドウGROUP BYは問題ではない」と記しています。
ジョブ全体に1つの値を設定するのは大雑把です。注文は1日覚えている必要があるのに、決済は2時間で済む結合があります。そのため、STATE_TTLヒントがあります。通常結合とGROUP BYにだけ使い、テーブル名またはエイリアスをキーとして指定します(エイリアスを付けたなら、必ずエイリアス)。
SET 'table.exec.state.ttl' = '30 min';
SELECT /*+ STATE_TTL('o' = '1d', 'p' = '2h') */ ...
FROM orders o JOIN payments p ON o.order_id = p.order_id;
-- 계획: leftState "1 d" · rightState "2 h" (잡 기본값보다 힌트가 우선)
このコードブロックの韓国語コメントは、計画にleftStateが「1 d」、rightStateが「2 h」と出て、ジョブのデフォルト値よりヒントが優先される、という意味です。
結合が連続する場合は、ルールがもう1つあります。ドキュメントは、ヒント3つが最初の結合の左側・右側と、2番目の結合の右側に解釈され、2番目の結合の左側(最初の結合の結果)は、ジョブの設定から来ると記しています。実際に、STATE_TTL('o'='1d','p'='2h','u'='7d')にジョブのデフォルト値30 minを指定して取り出すと、2番目の結合はleftState "30 min"・rightState "7 d"でした。その位置まで決めるには、最初の結合をビューに切り出して、別にヒントを指定する必要があります。
ステートバックエンドは、これらの状態をどこに置くかです。ドキュメントによると、何も指定しなければHashMapStateBackendで、状態をJavaヒープのオブジェクトとして置きます。EmbeddedRocksDBStateBackendは、TaskManagerのローカルディスクのRocksDBに、シリアライズされたバイトとして置きます。ディスクの容量ぶんまで入る代わりに、読み書きのたびにシリアライズのコストを払います。ドキュメントは、増分チェックポイントを提供するバックエンドとしてこちらを挙げ、リモートストレージに状態を置くForStバックエンドは、まだ実験段階だと記しています。ジョブごとにSET 'state.backend.type' = 'rocksdb'で切り替えられ、どのバックエンドで動いたかは、/jobs/<jid>/checkpoints/configのレスポンスのstate_backendの項目に残ります。同じGROUP BYをRocksDBで動かしても、結果は1文字も変わりません。
現場での姿
状態が増えているというアラートが来たら、まず計画を取り出します。どのオペレーターがTTL 0 msで何を保持しているかが、一目でわかります。たいてい原因は、何気なく書いた際限のないGROUP BYや、ディメンションテーブルとの通常結合です。ウィンドウ集計に変えられないか、テンポラル結合に変えられないかが、TTLより先に検討する解決策です。
TTLを設定すると決めたなら、ジョブ全体の値よりヒントを先に考えます。ジョブ全体で30分にすると、1日分の注文の状態まで30分で消えて、結果を間違えます。逆に、ヒントだけを信じて、連続した結合の2番目の左側の位置を忘れると、そこはジョブのデフォルト値(デフォルトは0 = 永遠)を受け取ります。計画JSONで数字を確認する習慣が、答えです。
バックエンドを変える理由は、たいていヒープが足りないからです。このPodのTaskManagerのタスクヒープは200MiBちょっとなので、大きな状態はhashmapに入りません。RocksDBに移すと、ディスクにあふれさせられますが、遅くなります。サイズと速度のトレードオフの判断です。
次のラボですること
3つのソースと、捨てるシンクを定義して、フィルターだけを行うクエリの計画で、状態がないことを確認します。際限のないGROUP BYのデフォルトのTTLを読み、ジョブ全体のTTLを設定して、変わった値を見ます。結合の両側に異なるTTLをヒントで指定し、連続した結合でヒントが届かない位置を探します。ウィンドウ集計の計画と結果を確認し、RocksDBバックエンドで同じ集計を動かしたあと、すべての数字を報告書にまとめます。