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

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

実行計画から状態と TTL を読む

TT Labで続きを見る

目標

オペレーターごとにどんな状態を保持していてTTLがいくつなのかを、COMPILE PLANのJSONで読み、ジョブ全体のTTLとオペレーターごとのヒントがどこに入るか、ウィンドウ集計とステートバックエンドが何を変えるかを確認します。

なぜ重要なのか

際限のないGROUP BYと通常結合は、デフォルトで状態を永遠に保持し、TTLで減らすと結果が間違うことがあります。TTLは「最後の更新からこれだけ経ったら、いつか削除する」なので、結果からは判定できませんが、計画は決定的に書き出されます。運用で状態の問題を扱う最初の道具が、まさにこの計画です。このラボの採点ツールは、クラスターに問い合わせず、皆さんが保存した計画JSON・RESTレスポンス・sql-clientの出力を読み、結果は元のCSVから直接計算して照合します。

ステップ

  1. 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に取り出してください。
  2. /root/flink/state/agg.sqlで、ユーザーごとのSUM(amount)を出す、際限のないGROUP BYの計画を、/root/flink/state/agg-plan.jsonに取り出してください(TTLの設定なしで)。
  3. /root/flink/state/agg-ttl.sqlで、SET 'table.exec.state.ttl' = '30 min';を入れた同じ計画を、/root/flink/state/agg-ttl-plan.jsonに取り出してください。
  4. /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に取り出してください。
  5. /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に取り出してください。
  6. /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に保存してください。
  7. /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に保存してください。
  8. /root/flink/state/report.jsonに、default_ttl・join_left_ttl・join_right_ttl・cascade_second_left_ttl・window_state_entries・state_backend・rocks_groupsを書いてください。

参考

フィルターだけのクエリには状態がない

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です。