TT Lab
Get started
Learn Learning paths Courses

Apache Flink — Running Streams on a Real Engine

Read the plan and predict the vertex count

Continue in TT Lab

Goal

Get the execution plan of the same GROUP BY query with EXPLAIN, read the operators and the place of the Exchange, and explain the number of vertices of the actual job by chaining and ship strategies. Confirm how mini-batch two-phase aggregation and distinct splitting appear in the plan.

Why it matters

The causes of a slow job — a filter that did not get pushed down to the source, an aggregation piled onto one key, a broken chain — look identical in a throughput graph and different in the plan. A tuning setting too has taken effect only if the plan changes. The grader of this lab does not ask the cluster anything. It reads only the EXPLAIN output text and REST responses you saved, and cross-checks the number of vertices against the ship strategies in the JSON execution plan.

Steps

  1. After flink-up, create the source clicks and the blackhole sink page_stats in /root/flink/plan/ddl.sql, and save the EXPLAIN PLAN FOR output of /root/flink/plan/explain.sql to /root/flink/plan/explain.out.
  2. Read the execution plan in explain.out and write exchange, above, below, and filter_pushed_into_scan to /root/flink/plan/shuffle.json.
  3. Save the output of /root/flink/plan/advice.sql, which uses EXPLAIN ESTIMATED_COST, PLAN_ADVICE, to /root/flink/plan/advice.out.
  4. Save the output of /root/flink/plan/json.sql, which uses EXPLAIN JSON_EXECUTION_PLAN INSERT INTO page_stats …, to /root/flink/plan/json.out.
  5. Run /root/flink/plan/chained.sql, which runs the same INSERT to the end with the job name flk-plan-chained, and save that job's /jobs/<jid> to /root/flink/plan/chained-job.json.
  6. Run /root/flink/plan/unchained.sql, which runs with chaining off and the job name flk-plan-unchained, and save /jobs/<jid> to /root/flink/plan/unchained-job.json.
  7. Save the output of /root/flink/plan/twophase.sql, which EXPLAINs the same SELECT after giving the three mini-batch settings, to /root/flink/plan/twophase.out.
  8. Save the output of /root/flink/plan/distinct.sql, which turns on distinct splitting (64 buckets) and EXPLAINs COUNT(DISTINCT user_id), to /root/flink/plan/distinct.out, and write /root/flink/plan/report.json collecting the numbers.

Notes

Get the three sections of EXPLAIN

After flink-up, in /root/flink/plan/ddl.sql write clicks, which reads plan_clicks.csv, and CREATE TABLE page_stats (page STRING, n_views BIGINT, revenue BIGINT) WITH ('connector' = 'blackhole'); in /root/flink/plan/explain.sql write EXPLAIN PLAN FOR plus the query from the Notes; and create /root/flink/plan/explain.out with sql-client.sh -i ddl.sql -f explain.sql > explain.out 2>&1.

EXPLAIN does not run the job but only returns the plan. The output must show the three sections: Abstract Syntax Tree, Optimized Physical Plan, and Optimized Execution Plan. In the initialization file (-i), put the two CREATE TABLE statements one per line.

Find the shuffle point

Read the execution plan section of explain.out and write to /root/flink/plan/shuffle.json exchange (the distribution value of the Exchange, for example hash[…]), above (the list of operator names above the Exchange), below (the list of operator names below it, from the top), and filter_pushed_into_scan (whether filter=[…] is attached to the source line, true or false).

An operator name is the word before the parenthesis (GroupAggregate, Calc, etc.). In the tree, the top is the sink side. If the WHERE was pushed down to the source, you see filter=[…] inside the TableSourceScan line. The grader compares with your explain.out.

Attach the cost estimate and advice

Run /root/flink/plan/advice.sql, which writes EXPLAIN ESTIMATED_COST, PLAN_ADVICE plus the same query, and save it to /root/flink/plan/advice.out. The cumulative cost and an advice[1]: [ADVICE] line must show up.

If you attach PLAN_ADVICE, the title of the physical plan section changes to 'With Advice' and an advice line is attached at the end. Read what it tells you to turn on — you will use it in step 7. The estimated cost is an assumed value for a source without statistics.

See the ship strategies in the JSON execution plan

Run /root/flink/plan/json.sql, which writes EXPLAIN JSON_EXECUTION_PLAN INSERT INTO page_stats plus the same query, and save it to /root/flink/plan/json.out. The nodes including the sink (Writer) and the ship_strategy must show up.

You have to EXPLAIN the INSERT, not the SELECT, to get the same graph as the job that will actually run. The nodes array under == Physical Execution Plan == is the operators, and the ship_strategy in predecessors is how data is received from the previous operator. Count how many times HASH appears.

Count the vertices of the job with chaining on

Run /root/flink/plan/chained.sql, which writes SET 'pipeline.name' = 'flk-plan-chained';, SET 'table.dml-sync' = 'true';, and INSERT INTO page_stats plus the same query, and then save that job's /jobs/<jid> to /root/flink/plan/chained-job.json. The number of vertices must be the number of non-FORWARD edges in json.out + 1.

If you turn on table.dml-sync, sql-client waits until the INSERT job finishes, so you can get the response of the finished (FINISHED) job. The operators joined with -> in a vertex name are a chain tied into one task.

Turn chaining off and count again

Run /root/flink/plan/unchained.sql, which is chained.sql with SET 'pipeline.operator-chaining.enabled' = 'false'; added and the job name changed to flk-plan-unchained, and then save /jobs/<jid> to /root/flink/plan/unchained-job.json. The number of vertices must equal the number of operators (nodes) in json.out.

If you turn chaining off, even operators connected by FORWARD each become a vertex (task). The result is the same and only the handoffs between tasks increase. If you do not change the name, you will confuse it with the job from the previous step in the job list.

Create two-phase aggregation with mini-batch

SET table.exec.mini-batch.enabled = true, table.exec.mini-batch.allow-latency = 1 s, and table.exec.mini-batch.size = 1000, and then run /root/flink/plan/twophase.sql, which writes EXPLAIN PLAN FOR plus the same query, and save it to /root/flink/plan/twophase.out. The execution plan must show GlobalGroupAggregate ← Exchange ← LocalGroupAggregate and MiniBatchAssigner below it.

The setting the advice in step 3 recommends. Two-phase aggregation arises only when mini-batch is on — pre-combine in each subtask before shuffling (Local), and combine after shuffling (Global). If even one of the three settings is missing, the plan does not change.

Distinct splitting and the report

SET table.optimizer.distinct-agg.split.enabled = true and table.optimizer.distinct-agg.split.bucket-num = 64, and save the output of /root/flink/plan/distinct.sql, which has EXPLAIN PLAN FOR SELECT page, COUNT(DISTINCT user_id) AS users FROM clicks GROUP BY page;, to /root/flink/plan/distinct.out. Then write to /root/flink/plan/report.json json_nodes, forward_edges, and hash_edges (json.out), chained_vertices and unchained_vertices (the two job responses), and distinct_buckets (the number of buckets in distinct.out).

If the split takes effect, the GroupAggregate becomes two layers, PARTIAL and FINAL, one Exchange appears between them and one below, and MOD(HASH_CODE(user_id), number of buckets) shows in the Calc below. All the numbers in the report can be counted from the files you saved — also check whether the chained vertex count is the HASH count + 1.