Apache Flink — Running Streams on a Real Engine
Reading execution plans — how many tasks does one SQL statement become?
In one line
EXPLAIN shows which operators an SQL statement is turned into, and the Exchange between them is the place where data is shuffled across the network. Operators that connect without shuffling (FORWARD) are tied into one task by chaining. So if you count the Exchanges in the plan, you can know in advance how many vertices the job will come up with, and you also check from the plan whether tunings such as two-phase aggregation or distinct splitting actually took effect.
Why this was needed
When you get a "the job is slow" report, you usually raise the parallelism first. But if the bottleneck is an aggregation piled onto one key, even if you raise the parallelism, the one subtask in charge of that key still receives everything. The filter may not be getting pushed down to the source so every row is being read, or chaining may be broken so handoffs between tasks have increased. These causes look identical in a throughput graph, and all look different in the plan.
The same goes for tuning options. Turning on table.exec.mini-batch.enabled does not by itself create a two-phase aggregation; the plan changes only when the optimizer judges the conditions to be met. The point of this module is: do not trust the configuration file, look at the plan.
How it works
The three sections of EXPLAIN. EXPLAIN PLAN FOR <질의> (where the placeholder stands for the query) produces three sections (EXPLAIN Statements in the official docs).
| Section | What |
|---|---|
== Abstract Syntax Tree == |
The logical plan that carries the SQL over as it is. LogicalFilter and LogicalAggregate |
== Optimized Physical Plan == |
The physical plan with rules applied. The filter moves down to the source and the aggregation becomes GroupAggregate |
== Optimized Execution Plan == |
The operator tree that will actually be built |
In the tree, the top is the sink side and the bottom is the source side. If you EXPLAIN this lab's query in streaming, you get GroupAggregate ← Exchange(distribution=[hash[page]]) ← Calc ← TableSourceScan, and filter=[>(amount, 0)] is attached to the source line — it means the WHERE was pushed down to the source. Exchange is the place where rows are shuffled by hash so that rows with the same key go to the same subtask, and the aggregation always stands above it.
There are details you can attach to EXPLAIN. ESTIMATED_COST attaches the estimated row count and cumulative cost for each node, CHANGELOG_MODE the kind of changelog (module 2), PLAN_ADVICE risk warnings and tuning advice, and JSON_EXECUTION_PLAN the operator graph as JSON. If you attach PLAN_ADVICE to a GROUP BY on this Pod, an [ADVICE] saying "try turning on local-global two-phase" appears. The estimated cost is an assumed value for a file source without statistics, so it must not be read as an absolute amount.
Chaining and vertices. The nodes in the JSON execution plan are individual operators, and between nodes the ship_strategy (FORWARD, HASH, etc.) is written. Operators connected by FORWARD with the same parallelism become a chain and run in one thread (module 1). So for a linear pipeline, the number of vertices = the number of non-FORWARD edges + 1. If you turn it off with pipeline.operator-chaining.enabled = false, each operator gets its own vertex — it is a switch to flip briefly when debugging to see per-operator metrics, not something to leave on in production.
Two-phase aggregation. Mini-batch is a device that gathers input briefly to reduce state access to once per key (Performance Tuning). When mini-batch is on, the optimizer can split an aggregation into a LocalGroupAggregate (before the shuffle, pre-combining in each subtask) and a GlobalGroupAggregate (combining after the shuffle). The docs liken this to the Combine + Reduce of MapReduce. Even if there is data piled onto one key, only the pre-combined accumulators cross the Exchange. On this Pod, when I gave just the three mini-batch settings (agg-phase-strategy default AUTO), MiniBatchAssigner and Local/Global appeared in the plan.
Distinct splitting. COUNT(DISTINCT user_id) does not shrink well even when pre-combined — because the accumulator ends up being a list of users. If you turn on table.optimizer.distinct-agg.split.enabled, the optimizer rewrites the query into two layers. The first layer adds MOD(HASH_CODE(user_id), 버킷 수) to the key and shuffles (PARTIAL) (where the last argument is the number of buckets), and the second layer shuffles again by the original key and adds up (FINAL). The default number of buckets is 1024.
What it looks like in the field
The scene you use most is "where does the filter apply". If the source supports filter pushdown, filter=[…] is attached to the TableSourceScan line. If it is not attached, every row leaves the source and is discarded only at the Calc.
The second is the hot key. If only one subtask is busy on the dashboard, look at the key of Exchange(distribution=[hash[…]]) in the plan. If the distribution of that key is skewed, the answer is not parallelism but two-phase aggregation or distinct splitting. After turning it on, always check with EXPLAIN whether Local/Global or PARTIAL/FINAL appeared.
The third is when the number of vertices differs from what you expected. The chain breaks between operators with different parallelism and at shuffle points. The vertex names in /jobs/<jid> list the operators joined with ->, so if you set that side by side with the ship strategies in the JSON execution plan, you can see right away where it broke.
What you will do in the next lab
You EXPLAIN the per-page aggregation of the click file to get the three sections, and transcribe what is above and below the Exchange into a JSON answer. After getting the EXPLAIN with cost and advice attached and the JSON execution plan, you run the same INSERT twice, with chaining on and off, and match the number of vertices with the ship strategies in the JSON. Finally, you get the plans of the mini-batch two-phase aggregation and the COUNT(DISTINCT) split and collect the numbers into a report.