Apache Flink — ストリームを本物のエンジンで動かす
実行計画から状態と TTL を読む
目標
オペレーターごとにどんな状態を保持していてTTLがいくつなのかを、COMPILE PLANのJSONで読み、ジョブ全体のTTLとオペレーターごとのヒントがどこに入るか、ウィンドウ集計とステートバックエンドが何を変えるかを確認します。
なぜ重要なのか
際限のないGROUP BYと通常結合は、デフォルトで状態を永遠に保持し、TTLで減らすと結果が間違うことがあります。TTLは「最後の更新からこれだけ経ったら、いつか削除する」なので、結果からは判定できませんが、計画は決定的に書き出されます。運用で状態の問題を扱う最初の道具が、まさにこの計画です。このラボの採点ツールは、クラスターに問い合わせず、皆さんが保存した計画JSON・RESTレスポンス・sql-clientの出力を読み、結果は元のCSVから直接計算して照合します。
ステップ
flink-upのあと、/root/flink/state/ddl.sqlに、orders・payments・usersと、blackholeシンクkv_sink(k STRING, v BIGINT)・win_sink(window_start, window_end, v BIGINT)を定義し、フィルターだけを行うINSERTの計画を、/root/flink/state/calc.sqlで、/root/flink/state/calc-plan.jsonに取り出してください。- /root/flink/state/agg.sqlで、ユーザーごとの
SUM(amount)を出す、際限のないGROUP BYの計画を、/root/flink/state/agg-plan.jsonに取り出してください(TTLの設定なしで)。 - /root/flink/state/agg-ttl.sqlで、
SET 'table.exec.state.ttl' = '30 min';を入れた同じ計画を、/root/flink/state/agg-ttl-plan.jsonに取り出してください。 - /root/flink/state/join-hint.sqlで、ジョブのデフォルト値30 minを置いたまま、orders o・payments pの結合に
STATE_TTL('o' = '1d', 'p' = '2h')を指定した計画を、/root/flink/state/join-hint-plan.jsonに取り出してください。 - /root/flink/state/cascade.sqlで、ジョブのデフォルト値30 minに、orders o ⋈ payments p ⋈ users uの結合と
STATE_TTL('o' = '1d', 'p' = '2h', 'u' = '7d')を指定した計画を、/root/flink/state/cascade-plan.jsonに取り出してください。 - /root/flink/state/window.sqlの1つのファイルで、ジョブのデフォルト値30 minを置いたまま、10分のTUMBLEウィンドウ集計の計画を、/root/flink/state/window-plan.jsonに取り出し、同じウィンドウの
window_start・window_end・orders・amountをストリーミングで動かして、出力を、/root/flink/state/window.outに保存してください。 - /root/flink/state/rocks.sqlに、
SET 'state.backend.type' = 'rocksdb'・ジョブ名flk-state-rocksで、ユーザーごとのorders・totalを出すクエリを書いて実行し、出力を、/root/flink/state/rocks.outに、そのジョブの/jobs/<jid>を、/root/flink/state/rocks-job.jsonに、/jobs/<jid>/checkpoints/configを、/root/flink/state/rocks-ckpt.jsonに保存してください。 - /root/flink/state/report.jsonに、
default_ttl・join_left_ttl・join_right_ttl・cascade_second_left_ttl・window_state_entries・state_backend・rocks_groupsを書いてください。
参考
- 元データ(ヘッダーなしのCSV、時刻は秒単位、tsは昇順):
state_orders.csv=order_id, user_id, amount, ts・state_payments.csv=pay_id, order_id, pay_method, ts・state_users.csv=user_id, region。すべて/opt/lab/fixtures/data/にあります。orders・paymentsのtsにウォーターマークを設定してください。 - 計画の取り出し:
COMPILE PLAN 'file:///root/flink/state/이름.json' FOR INSERT INTO 싱크 SELECT ...;(プレースホルダーは順に、ファイル名とシンクです)。ジョブは動かしません。同じパスにファイルがあると、上書きせずにエラーになるので、取り出し直すときは先に削除します。 - 計画の確認:
jq -c '.nodes[] | {id, type, state}' 파일.json(プレースホルダーはファイル名です)。ノードidの順序が、下(ソース)から上(シンク)への順序です。 - よくある間違い: ドキュメントどおりテーブルにエイリアスを付けたなら、ヒントのキーもエイリアスでなければなりません。
methodはSQLの予約語なので、列名に使うとCREATEが失敗します(pay_method)。 - 公式ドキュメント: Hints — STATE_TTL・Configuration — table.exec.state.ttl・State Backends・Group Aggregation・Determinism
フィルターだけのクエリには状態がない
flink-upのあと、/root/flink/state/ddl.sqlに、orders・payments・usersと、blackholeシンクkv_sink・win_sinkを定義してください。/root/flink/state/calc.sqlに、amount > 1000の注文のorder_idとamount * 2をkv_sinkに入れるINSERTを、COMPILE PLAN 'file:///root/flink/state/calc-plan.json'で取り出す文を書き、sql-client.sh -i ddl.sql -f calc.sqlで実行してください。
COMPILE PLANはINSERT文を受け取るので、受け取るシンクが必要です。blackholeコネクターは、受け取ったものを捨てます。作成したJSONをjqで開いて、nodesのtypeとstateを見てください。行を1つ見てすぐ出力するオペレーターには、記憶するものがありません。
際限のないGROUP BY: デフォルトは永遠
/root/flink/state/agg.sqlで、ユーザーごとのSUM(amount)をkv_sinkに入れるINSERTの計画を、/root/flink/state/agg-plan.jsonに取り出してください。TTLは設定しません。
集計ノードのstateのリストに何があり、ttlがいくつかを見てください。0 msは、削除しないという意味です。SUMの結果の型がシンクの列(BIGINT)と違う場合は、CASTで合わせます。
ジョブ全体のTTLを設定する
/root/flink/state/agg-ttl.sqlの先頭にSET 'table.exec.state.ttl' = '30 min';を置き、ステップ2と同じGROUP BYの計画を、/root/flink/state/agg-ttl-plan.jsonに取り出してください。
SETは、そのあとの文から適用されます。計画のgroupAggregateStateのttlがどう変わるかを見てください。この値は「最後の更新から最低でもこれだけは保持」という意味なので、結果からは、いつ削除されたかがわかりません。
結合の両側に異なるTTL: STATE_TTLヒント
/root/flink/state/join-hint.sqlに、SET 'table.exec.state.ttl' = '30 min';を置いたまま、orders oとpayments pをorder_idで通常結合してkv_sinkに入れるINSERTに、/*+ STATE_TTL('o' = '1d', 'p' = '2h') */を指定して、計画を、/root/flink/state/join-hint-plan.jsonに取り出してください。
ヒントは、SELECTの直後に書きます。テーブルにエイリアスを付けたなら、ヒントのキーもエイリアスでなければなりません。計画の結合ノードで、leftState・rightStateのttlが、ジョブのデフォルト値と違って出るかを見てください。
連続した結合: ヒントが届かない位置
/root/flink/state/cascade.sqlに、ジョブのデフォルト値30 minを置き、orders o ⋈ payments p(order_id) ⋈ users u(user_id)をkv_sinkに入れるINSERTに/*+ STATE_TTL('o' = '1d', 'p' = '2h', 'u' = '7d') */を指定して、計画を、/root/flink/state/cascade-plan.jsonに取り出してください。
結合ノードが2つ出ます。idが小さいほうが、先に(下で)行われる結合です。ヒント3つが4つの位置(最初の結合の左側・右側、2番目の結合の左側・右側)のどこに付き、残った1つの位置は何を受け取るかを見てください。
ウィンドウ集計: TTLなしで自分で空にする
/root/flink/state/window.sqlの1つのファイルに、SET 'table.exec.state.ttl' = '30 min';を置き、(1)10分のTUMBLEウィンドウのSUM(amount)をwin_sinkに入れる計画を、/root/flink/state/window-plan.jsonに取り出し、(2)同じウィンドウのwindow_start・window_end・orders(COUNT)・amount(SUM)をストリーミングのSELECTで動かしてください。出力は、/root/flink/state/window.outに保存します。
ウィンドウTVFは、FROM TABLE(TUMBLE(TABLE orders, DESCRIPTOR(ts), INTERVAL '10' MINUTE))で、GROUP BY window_start, window_endでまとめます。TTLを設定したのに、ウィンドウ集計のノードにstateの項目があるかを見てください。結果のopがすべて+Iなのも、同じ理由です。ウィンドウは、ウォーターマークが終わりを過ぎたときに1回出して、空にします。
RocksDBバックエンドで同じ集計を動かす
/root/flink/state/rocks.sqlに、SET 'state.backend.type' = 'rocksdb';・SET 'pipeline.name' = 'flk-state-rocks';を置き、ユーザーごとのorders(COUNT)・total(SUM(amount))をストリーミングで出してください。出力は、/root/flink/state/rocks.outに、そのジョブの/jobs/<jid>は、/root/flink/state/rocks-job.jsonに、/jobs/<jid>/checkpoints/configは、/root/flink/state/rocks-ckpt.jsonに保存します。
ジョブidは、/jobs/overviewから名前で探します。checkpoints/configのレスポンスのstate_backendの項目に、このジョブが使ったバックエンドのクラス名が書かれます。デフォルトのバックエンドで動かしたジョブなら、別の名前が出ます。結果は、バックエンドと関係なく同じでなければなりません。
報告書: 状態の地図を数字で
/root/flink/state/report.jsonに、default_ttl(agg-planのgroupAggregateStateのttl)、join_left_ttl・join_right_ttl(join-hint-plan)、cascade_second_left_ttl(cascade-planの2番目の結合のleftStateのttl)、window_state_entries(window-planのstateの項目数、整数)、state_backend(rocks-ckptの値)、rocks_groups(rocks.outの最終結果の行数 = 状態に残ったキー数、整数)を書いてください。
ttlの値は、計画に書かれた文字列(例: "2 h")をそのまま写せばよいです。jqで、typeがstream-exec-joinで始まるノードをid順に並べれば、2つの結合を区別できます。stateがないノードは、その項目がnullです。