Apache Flink — ストリームを本物のエンジンで動かす
計画を読んで頂点数を当てる
目標
同じGROUP BYクエリの実行計画をEXPLAINで取得して、オペレーターとExchangeの位置を読み、実際のジョブの頂点数を、チェイニングと伝達方式で説明します。mini-batchの2フェーズ集計とdistinct分割が、計画にどう現れるかを確認します。
なぜ重要なのか
遅いジョブの原因(ソースまで下がらなかったフィルター、1つのキーに集中した集計、途切れたチェーン)は、スループットのグラフでは同じに見え、計画では違って見えます。チューニング設定も、計画が変わってはじめて効いたことになります。このラボの採点ツールは、クラスターに問い合わせません。皆さんが保存したEXPLAIN出力の原文とRESTレスポンスだけを読み、頂点数は、JSON実行計画の伝達方式とクロスチェックして検証します。
ステップ
flink-upのあと、/root/flink/plan/ddl.sqlに、ソースclicksとblackholeシンクpage_statsを作成し、/root/flink/plan/explain.sqlのEXPLAIN PLAN FORの出力を、/root/flink/plan/explain.outに保存してください。- explain.outの実行計画を読んで、/root/flink/plan/shuffle.jsonに、
exchange・above・below・filter_pushed_into_scanを書いてください。 - /root/flink/plan/advice.sqlに
EXPLAIN ESTIMATED_COST, PLAN_ADVICEを書いて実行し、出力を、/root/flink/plan/advice.outに保存してください。 - /root/flink/plan/json.sqlに
EXPLAIN JSON_EXECUTION_PLAN INSERT INTO page_stats …を書いて実行し、出力を、/root/flink/plan/json.outに保存してください。 - /root/flink/plan/chained.sqlに、ジョブ名
flk-plan-chainedで同じINSERTを最後まで動かすSQLを書いて実行し、そのジョブの/jobs/<jid>を、/root/flink/plan/chained-job.jsonに保存してください。 - /root/flink/plan/unchained.sqlに、チェイニングを無効にして、ジョブ名
flk-plan-unchainedで動かすSQLを書いて実行し、/jobs/<jid>を、/root/flink/plan/unchained-job.jsonに保存してください。 - /root/flink/plan/twophase.sqlに、mini-batchの3つの設定を指定したあと、同じSELECTをEXPLAINするSQLを書いて実行し、出力を、/root/flink/plan/twophase.outに保存してください。
- /root/flink/plan/distinct.sqlに、distinct分割(バケット64)を有効にして、
COUNT(DISTINCT user_id)をEXPLAINするSQLを書いて実行し、出力を、/root/flink/plan/distinct.outに保存してください。そして、数字を集めた報告書を、/root/flink/plan/report.jsonに書いてください。
参考
- 元データ
/opt/lab/fixtures/data/plan_clicks.csvの列:click_id BIGINT, user_id STRING, page STRING, amount INT, click_time TIMESTAMP(3)(ヘッダーなしのCSV)。 - クエリは、すべてのステップで同じです:
SELECT page, COUNT(*) AS n_views, SUM(amount) AS revenue FROM clicks WHERE amount > 0 GROUP BY page。viewsは予約語なので、別名には使えません。 - テーブル定義を毎回繰り返さないために、
sql-client.sh -i ddl.sql -f 파일.sql > 파일.out 2>&1(プレースホルダーはファイル名です)を使います(-iは、先に実行される初期化ファイル)。sql-clientは、1行に1つの文だけを受け付けます。 - EXPLAINの出力は、表の1つのセルに複数行で出力されます。実行計画は、上がシンク側、下がソース側です。
- よくある間違い:
SET 'execution.runtime-mode' = 'batch'で動かすと、集計オペレーターの名前が変わります。このラボは、デフォルト(ストリーミング)で行います。INSERTが終わる前に/jobs/<jid>を取得すると、RUNNINGです(table.dml-sync)。 - 公式ドキュメント: EXPLAIN・Performance Tuning・Tableの設定・Flink Architecture・REST API
EXPLAINの3つのセクションを取得する
flink-upのあと、/root/flink/plan/ddl.sqlに、plan_clicks.csvを読むclicksとCREATE TABLE page_stats (page STRING, n_views BIGINT, revenue BIGINT) WITH ('connector' = 'blackhole')を書き、/root/flink/plan/explain.sqlにEXPLAIN PLAN FORと参考のクエリを書いて、sql-client.sh -i ddl.sql -f explain.sql > explain.out 2>&1で、/root/flink/plan/explain.outを作成してください。
EXPLAINは、ジョブを動かさず、計画だけを返します。出力にAbstract Syntax Tree・Optimized Physical Plan・Optimized Execution Planの3つのセクションが見える必要があります。初期化ファイル(-i)には、CREATE TABLEの2つの文を、1行に1つずつ置きます。
シャッフルの位置を見つける
explain.outの実行計画セクションを読んで、/root/flink/plan/shuffle.jsonに、exchange(Exchangeのdistributionの値、例: hash[…])、above(Exchangeより上のオペレーター名の一覧)、below(下のオペレーター名の一覧、上から)、filter_pushed_into_scan(ソースの行にfilter=[…]が付いているか、true/false)を書いてください。
オペレーター名は、括弧の前の単語です(GroupAggregate、Calcなど)。ツリーは上がシンク側です。WHEREがソースまで下がっていれば、TableSourceScanの行の中にfilter=[…]が見えます。採点ツールは、皆さんのexplain.outと照合します。
コスト推定とアドバイスを付ける
/root/flink/plan/advice.sqlに、EXPLAIN ESTIMATED_COST, PLAN_ADVICEと同じクエリを書いて実行し、出力を、/root/flink/plan/advice.outに保存してください。コスト(cumulative cost)とadvice[1]: [ADVICE]の行が見える必要があります。
PLAN_ADVICEを付けると、物理計画セクションの見出しが「With Advice」に変わり、末尾にアドバイスの行が付きます。何を有効にしてみるよう勧めているかを読んでおいてください。ステップ7で使います。推定コストは、統計のないソースに対する仮定値です。
JSON実行計画で伝達方式を見る
/root/flink/plan/json.sqlに、EXPLAIN JSON_EXECUTION_PLAN INSERT INTO page_statsと同じクエリを書いて実行し、出力を、/root/flink/plan/json.outに保存してください。シンク(Writer)まで含めたノードとship_strategyが見える必要があります。
SELECTではなくINSERTをEXPLAINしてはじめて、実際に動かすジョブと同じグラフが出ます。== Physical Execution Plan ==の下のnodes配列がオペレーターで、predecessorsのship_strategyが、前のオペレーターからデータを受け取る方式です。HASHが何回出るかを数えておいてください。
チェイニングを有効にしたジョブの頂点を数える
/root/flink/plan/chained.sqlに、SET 'pipeline.name' = 'flk-plan-chained';・SET 'table.dml-sync' = 'true';・INSERT INTO page_statsと同じクエリを書いて実行したあと、そのジョブの/jobs/<jid>を、/root/flink/plan/chained-job.jsonに保存してください。頂点数は、json.outのFORWARD以外の伝達数 + 1である必要があります。
table.dml-syncを有効にすると、sql-clientがINSERTジョブが終わるまで待つので、終わった(FINISHED)ジョブのレスポンスを取得できます。頂点名に->でつながれたオペレーターが、1つのタスクにまとめられたチェーンです。
チェイニングを無効にして、もう一度数える
/root/flink/plan/unchained.sqlに、chained.sqlにSET 'pipeline.operator-chaining.enabled' = 'false';を加えて、ジョブ名をflk-plan-unchainedに変えたSQLを書いて実行し、/jobs/<jid>を、/root/flink/plan/unchained-job.jsonに保存してください。頂点数が、json.outのオペレーター(ノード)数と同じである必要があります。
チェイニングを無効にすると、FORWARDでつながったオペレーターも、それぞれ頂点(タスク)になります。結果は同じで、タスク間の受け渡しだけが増えます。名前を変えないと、ジョブ一覧で、前のステップのジョブと区別がつかなくなります。
mini-batchで2フェーズ集計を作る
/root/flink/plan/twophase.sqlに、table.exec.mini-batch.enabled = true、table.exec.mini-batch.allow-latency = 1 s、table.exec.mini-batch.size = 1000をSETしたあと、EXPLAIN PLAN FORと同じクエリを書いて実行し、出力を、/root/flink/plan/twophase.outに保存してください。実行計画にGlobalGroupAggregate ← Exchange ← LocalGroupAggregateと、その下のMiniBatchAssignerが見える必要があります。
ステップ3のアドバイスが勧めた設定です。2フェーズ集計は、mini-batchが有効になってはじめてできます。シャッフルの前にサブタスクごとに事前に合算し(Local)、シャッフルのあとで合算します(Global)。3つの設定のうち1つでも欠けると、計画は変わりません。
distinct分割と報告書
/root/flink/plan/distinct.sqlに、table.optimizer.distinct-agg.split.enabled = true、table.optimizer.distinct-agg.split.bucket-num = 64をSETして、EXPLAIN PLAN FOR SELECT page, COUNT(DISTINCT user_id) AS users FROM clicks GROUP BY page;を書いて実行し、出力を、/root/flink/plan/distinct.outに保存してください。そして、/root/flink/plan/report.jsonに、json_nodes・forward_edges・hash_edges(json.out)、chained_vertices・unchained_vertices(2つのジョブのレスポンス)、distinct_buckets(distinct.outのバケット数)を書いてください。
分割が効けば、GroupAggregateがPARTIALとFINALの2層になり、その間と下にExchangeが1つずつでき、下のCalcにMOD(HASH_CODE(user_id), バケット数)が見えます。報告書の数字は、すべて皆さんが保存したファイルから数えられます。チェイニングの頂点数がHASHの数 + 1かも確認してみてください。