Apache Flink — ストリームを本物のエンジンで動かす
実行計画を読む — SQL 一文はいくつのタスクになるか
一言でいうと
EXPLAINは、SQLがどんなオペレーターに変わるかを見せてくれます。その間にあるExchangeが、データをネットワークで混ぜる位置です。混ぜずにつながる(FORWARD)オペレーターは、チェイニングで1つのタスクにまとめられます。そのため、計画のExchangeを数えれば、ジョブが頂点いくつで起動するかをあらかじめ知ることができ、2フェーズ集計やdistinct分割のようなチューニングが実際に効いたかも、計画で確認できます。
なぜ必要なのか
「ジョブが遅い」という報告を受けると、たいてい並列度から上げます。しかし、ボトルネックが1つのキーに集中した集計なら、並列度を上げても、そのキーを担当するサブタスク1つが、相変わらずすべてを受け取ります。フィルターがソースまで下がらず、すべての行を読んでいるのかもしれませんし、チェイニングが途切れて、タスク間の受け渡しが増えているのかもしれません。これらの原因は、スループットのグラフでは同じに見え、計画ではすべて違って見えます。
チューニングオプションも同じです。table.exec.mini-batch.enabledを有効にしたからといって2フェーズ集計ができるわけではなく、オプティマイザーが条件を満たすと判断したときだけ、計画が変わります。設定ファイルを信じず、計画を見よ、というのが、このモジュールの要点です。
どう動くのか
EXPLAINの3つのセクションは、次のとおりです。EXPLAIN PLAN FOR <질의>(プレースホルダーはクエリです)は、3つのセクションを出力します(公式ドキュメントのEXPLAIN Statements)。
| セクション | 内容 |
|---|---|
== Abstract Syntax Tree == |
SQLをそのまま写した論理計画。LogicalFilter・LogicalAggregate |
== Optimized Physical Plan == |
ルールを適用した物理計画。フィルターがソースまで下がり、集計がGroupAggregateになる |
== Optimized Execution Plan == |
実際に作られるオペレーターのツリー |
ツリーは、上がシンク側、下がソース側です。このラボのクエリをストリーミングでEXPLAINすると、GroupAggregate ← Exchange(distribution=[hash[page]]) ← Calc ← TableSourceScanが出て、ソースの行にfilter=[>(amount, 0)]が付きます。WHEREがソースまで下がったという意味です。Exchangeは、同じキーの行を同じサブタスクに送るために、ハッシュで混ぜる位置で、集計は必ずその上に立ちます。
EXPLAINには、付けられる詳細があります。ESTIMATED_COSTはノードごとの推定行数と累積コストを、CHANGELOG_MODEはチェンジログの種類(モジュール2)を、PLAN_ADVICEはリスクの警告とチューニングのアドバイスを、JSON_EXECUTION_PLANはオペレーターのグラフをJSONで付けます。このPodでGROUP BYにPLAN_ADVICEを付けると、「local-global two-phaseを有効にしてみてください」という[ADVICE]が出ます。推定コストは、統計のないファイルソースに対する仮定値なので、絶対量として読んではいけません。
チェイニングと頂点は、次のとおりです。JSON実行計画のノードは、オペレーター1つ1つで、ノードの間にship_strategy(FORWARD・HASHなど)が書かれます。FORWARDでつながり、並列度が同じオペレーターは、チェーンになって1つのスレッドで動きます(モジュール1)。そのため、線形のパイプラインなら、頂点数 = FORWARD以外の伝達数 + 1です。pipeline.operator-chaining.enabled = falseで無効にすると、オペレーターごとに頂点ができます。デバッグのときにオペレーター別の指標を見るために、一時的に切るスイッチであって、運用で無効のままにしておくものではありません。
2フェーズ集計は、次のとおりです。mini-batchは、入力を少しのあいだ集めて、キーごとの状態アクセスを1回に減らす仕組みです(Performance Tuning)。mini-batchが有効になると、オプティマイザーは集計を、LocalGroupAggregate(シャッフルの前に、各サブタスクで事前に合算)とGlobalGroupAggregate(シャッフルのあとで合算)に分けられます。ドキュメントは、これをMapReduceのCombine + Reduceになぞらえています。1つのキーにデータが集中していても、Exchangeを越えるのは、事前に合算したアキュムレーターだけです。このPodでは、mini-batchの3つの設定だけを指定すると(agg-phase-strategyのデフォルトはAUTO)、計画にMiniBatchAssignerとLocal/Globalが現れました。
distinct分割は、次のとおりです。COUNT(DISTINCT user_id)は、事前に合算してもあまり減りません。アキュムレーターが、結局はユーザーのリストだからです。table.optimizer.distinct-agg.split.enabledを有効にすると、オプティマイザーがクエリを2層に書き換えます。1層目は、MOD(HASH_CODE(user_id), 버킷 수)(プレースホルダーはバケット数です)をキーに加えてシャッフルし(PARTIAL)、2層目が元のキーでもう一度シャッフルして合算します(FINAL)。バケット数のデフォルトは1024です。
現場での姿
最もよく使う場面は、「フィルターがどこでかかるか」です。ソースがフィルターのプッシュダウンをサポートしていれば、TableSourceScanの行にfilter=[…]が付きます。付いていなければ、すべての行がソースを出て、Calcではじめて捨てられます。
2つ目は、ホットキーです。ダッシュボードで、サブタスク1つだけが忙しいなら、計画のExchange(distribution=[hash[…]])のキーを見ます。そのキーの分布が偏っているなら、答えは並列度ではなく、2フェーズ集計かdistinct分割です。有効にしたあとは、必ずEXPLAINで、Local/GlobalまたはPARTIAL/FINALができたかを確認します。
3つ目は、頂点数が予想と違うときです。並列度が異なるオペレーターの間や、シャッフルの位置の間では、チェーンが途切れます。/jobs/<jid>の頂点名に->でつながれたオペレーターの一覧が出るので、JSON実行計画の伝達方式と並べて見れば、どこで途切れたかがすぐわかります。
次のラボですること
クリックファイルのページ別集計をEXPLAINして3つのセクションを取得し、Exchangeの上下をJSONの答案に書き写します。コスト・アドバイスを付けたEXPLAINとJSON実行計画を取得したあと、同じINSERTをチェイニングを有効にした場合と無効にした場合の2回動かして、頂点数をJSONの伝達方式と合わせます。最後に、mini-batchの2フェーズ集計とCOUNT(DISTINCT)分割の計画を取得して、数字を報告書にまとめます。