Apache Flink — Running Streams on a Real Engine
Read the plan and predict the vertex count
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
- After
flink-up, create the sourceclicksand the blackhole sinkpage_statsin /root/flink/plan/ddl.sql, and save theEXPLAIN PLAN FORoutput of /root/flink/plan/explain.sql to /root/flink/plan/explain.out. - Read the execution plan in explain.out and write
exchange,above,below, andfilter_pushed_into_scanto /root/flink/plan/shuffle.json. - Save the output of /root/flink/plan/advice.sql, which uses
EXPLAIN ESTIMATED_COST, PLAN_ADVICE, to /root/flink/plan/advice.out. - 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. - 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. - 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. - 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.
- 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
- Columns of the source
/opt/lab/fixtures/data/plan_clicks.csv:click_id BIGINT, user_id STRING, page STRING, amount INT, click_time TIMESTAMP(3)(a CSV with no header). - The query is the same in every step:
SELECT page, COUNT(*) AS n_views, SUM(amount) AS revenue FROM clicks WHERE amount > 0 GROUP BY page.viewsis a reserved word and cannot be used as an alias. - To avoid repeating the table definitions every time, use
sql-client.sh -i ddl.sql -f 파일.sql > 파일.out 2>&1(the placeholders stand for the file names; -i is an initialization file that runs first). sql-client accepts only one statement per line. - EXPLAIN output is printed on several lines in one table cell. In the execution plan, the top is the sink side and the bottom is the source side.
- A common mistake: if you run with
SET 'execution.runtime-mode' = 'batch', the names of the aggregation operators change — this lab uses the default (streaming). If you get/jobs/<jid>before the INSERT ends, it is RUNNING (table.dml-sync). - Official docs: EXPLAIN · Performance Tuning · Table configuration · Flink Architecture · REST API
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.