Shows the execution plan of a statement.
Syntax:
EXPLAIN [AST | SYNTAX | QUERY TREE | PLAN | PIPELINE | ANALYZE | ESTIMATE | TABLE OVERRIDE | WHATIF] [setting = value, ...]
[
SELECT ... |
tableFunction(...) [COLUMNS (...)] [ORDER BY ...] [PARTITION BY ...] [PRIMARY KEY] [SAMPLE BY ...] [TTL ...]
]
[FORMAT ...]Example:
EXPLAIN SELECT sum(number) FROM numbers(10) UNION ALL SELECT sum(number) FROM numbers(10) ORDER BY sum(number) ASC FORMAT TSV;Output: sum(number)
Union
├──Aggregating
│ │ Keys:
│ │ Aggregates: sum(number)
│ │ Skip merging: 0
│ └──ReadFromSystemNumbers
│ Output: number
└──Sorting (Sorting for ORDER BY)
│ Sort description: sum(number) ASC
└──Aggregating
│ Keys:
│ Aggregates: sum(number)
│ Skip merging: 0
└──ReadFromSystemNumbers
Output: numberEXPLAIN Types
AST— Abstract syntax tree.SYNTAX— Query text after AST-level optimizations.QUERY TREE— Query tree after Query Tree level optimizations.PLAN— Query execution plan.PIPELINE— Query execution pipeline.ANALYZE— Executes the query and annotates the execution plan with measured runtime metrics.ESTIMATE— Estimated number of rows, marks and parts to be read from the tables while processing the query.TABLE OVERRIDE— Validated result of a table override on a table-function schema.
EXPLAIN AST
Dump query AST. Supports all types of queries, not only SELECT.
Settings:
graph– Prints AST as a graph described in the DOT graph description language. Default: 0.
Examples:
EXPLAIN AST SELECT 1;SelectWithUnionQuery (children 1)
ExpressionList (children 1)
SelectQuery (children 1)
ExpressionList (children 1)
Literal UInt64_1EXPLAIN AST ALTER TABLE t1 DELETE WHERE date = today(); explain
AlterQuery t1 (children 1)
ExpressionList (children 1)
AlterCommand 27 (children 1)
Function equals (children 1)
ExpressionList (children 2)
Identifier date
Function today (children 1)
ExpressionListEXPLAIN SYNTAX
Shows the Abstract Syntax Tree (AST) of a query after syntax analysis.
It’s done by parsing the query, constructing query AST and query tree, optionally running query analyzer and optimization passes, and then converting the query tree back to the query AST.
Settings:
oneline– Print the query in one line. Default:0.run_query_tree_passes– Run query tree passes before dumping the query tree. Default:0.query_tree_passes– Ifrun_query_tree_passesis set, specifies how many passes to run. Without specifyingquery_tree_passesit runs all the passes.single_record– Return the reformatted query as a single multi-line record instead of one record per line. Default:1(controlled by theexplain_syntax_single_recordsetting). Set to0to restore the historical one-record-per-line output, or setexplain_syntax_single_record = 0(globally or in per-querySETTINGS), or setcompatibilityto any version older than26.8.
Examples:
EXPLAIN SYNTAX SELECT * FROM system.numbers AS a, system.numbers AS b, system.numbers AS c WHERE a.number = b.number AND b.number = c.number;SELECT *
FROM system.numbers AS a, system.numbers AS b, system.numbers AS c
WHERE (a.number = b.number) AND (b.number = c.number)With run_query_tree_passes:
EXPLAIN SYNTAX run_query_tree_passes = 1 SELECT * FROM system.numbers AS a, system.numbers AS b, system.numbers AS c WHERE a.number = b.number AND b.number = c.number;SELECT
__table1.number AS `a.number`,
__table2.number AS `b.number`,
__table3.number AS `c.number`
FROM system.numbers AS __table1
ALL INNER JOIN system.numbers AS __table2 ON __table1.number = __table2.number
ALL INNER JOIN system.numbers AS __table3 ON __table2.number = __table3.numberEXPLAIN QUERY TREE
Settings:
run_passes— Run all query tree passes before dumping the query tree. Default:1.dump_passes— Dump information about used passes before dumping the query tree. Default:0.passes— Specifies how many passes to run. If set to-1, runs all the passes. Default:-1.dump_tree— Display the query tree. Default:1.dump_ast— Display the query AST generated from the query tree. Default:0.
Example:
EXPLAIN QUERY TREE SELECT id, value FROM test_table;QUERY id: 0
PROJECTION COLUMNS
id UInt64
value String
PROJECTION
LIST id: 1, nodes: 2
COLUMN id: 2, column_name: id, result_type: UInt64, source_id: 3
COLUMN id: 4, column_name: value, result_type: String, source_id: 3
JOIN TREE
TABLE id: 3, table_name: default.test_tableEXPLAIN PLAN
Dump query plan steps.
Settings:
optimize— Controls whether query plan optimizations are applied before displaying the plan. Default: 1.header— Prints output header for step. Default: 0.description— Prints step description. Default: 1.indexes— Shows used indexes, the number of filtered parts and the number of filtered granules for every index applied. Default: 0. Supported for MergeTree tables. Starting from ClickHouse >= v25.9, this statement only shows reasonable output when used withSETTINGS use_query_condition_cache = 0, use_skip_indexes_on_data_read = 0.projections— Shows all analyzed projections and their effect on part-level filtering based on projection primary key conditions. For each projection, this section includes statistics such as the number of parts, rows, marks, and ranges that were evaluated using the projection’s primary key. It also shows how many data parts were skipped due to this filtering, without reading from the projection itself. Whether a projection was actually used for reading or only analyzed for filtering can be determined by thedescriptionfield. Default: 0. Supported for MergeTree tables.actions— Prints detailed information about step actions. Default: 1.sorting— Prints the sort description for each plan step that produces sorted output. Default: 0.keep_logical_steps— Keeps logical plan steps for joins instead of converting them to physical join implementations. Default: 0.json— Prints query plan steps as a row in JSON format. Default: 0. It is recommended to use TabSeparatedRaw (TSVRaw) format to avoid unnecessary escaping.input_headers— Prints input headers for step. Default: 0. Mostly useful only for developers to debug issues related to input-output header mismatch.column_structure— Prints also the structure of columns in headers on top of their name and type. Default: 0. Mostly useful only for developers to debug issues related to input-output header mismatch.distributed— Shows query plans executed on remote nodes for distributed tables or parallel replicas. Not supported together withjson. Default: 0.compact— When enabled, hides expression steps and detailed action info (inputs, functions, aliases, and output positions) from the plan. Only has an effect whenactions = 1. Default: 1.pretty— Prints the plan tree using line-drawing characters (├──, └──, │) instead of indentation to visualize the hierarchy. Also formats join step properties inline. Default: 1.
Example:
EXPLAIN SELECT sum(number) FROM numbers(10) GROUP BY number % 4 LIMIT 1;Output: sum(number)
Limit (preliminary LIMIT)
│ Limit 1
│ Offset 0
└──Aggregating
│ Keys: number MOD 4
│ Aggregates: sum(number)
│ Skip merging: 0
└──ReadFromSystemNumbers
Output: numberWhen json = 1, the query plan is represented in JSON format. Every node is a dictionary that always has the keys Node Type, Node Id, and Plans. Node Type is a string with the step name, and Node Id is a unique step identifier (the step name with a numeric suffix, e.g. Union_10). Plans is an array with child step descriptions. Other optional keys may be added depending on node type and settings.
Example:
EXPLAIN json = 1, description = 0 SELECT 1 UNION ALL SELECT 2 FORMAT TSVRaw;[
{
"Plan": {
"Node Type": "Union",
"Node Id": "Union_10",
"Plans": [
{
"Node Type": "Expression",
"Node Id": "Expression_13",
"Plans": [
{
"Node Type": "ReadFromStorage",
"Node Id": "ReadFromStorage_0"
}
]
},
{
"Node Type": "Expression",
"Node Id": "Expression_16",
"Plans": [
{
"Node Type": "ReadFromStorage",
"Node Id": "ReadFromStorage_4"
}
]
}
]
}
}
]With description = 1, the Description key is added to the step:
{
"Node Type": "ReadFromStorage",
"Description": "SystemOne"
}With header = 1, the Header key is added to the step as an array of columns.
Example:
EXPLAIN json = 1, description = 0, header = 1 SELECT 1, 2 + dummy;[
{
"Plan": {
"Node Type": "Expression",
"Node Id": "Expression_5",
"Header": [
{
"Name": "1",
"Type": "UInt8"
},
{
"Name": "plus(2, dummy)",
"Type": "UInt16"
}
],
"Plans": [
{
"Node Type": "ReadFromStorage",
"Node Id": "ReadFromStorage_0",
"Header": [
{
"Name": "dummy",
"Type": "UInt8"
}
]
}
]
}
}
]With indexes = 1, the Indexes key is added. It contains an array of used indexes. Each index is described as JSON with Type key (a string Partition Min-Max, Partition, Statistics, PrimaryKey or Skip) and optional keys:
The Statistics index uses per-part column statistics (min/max values, and the number of NULL values for Nullable columns) to skip parts that cannot match the query filter.
Name— The index name (currently only used forSkipindexes).Keys— The array of columns used by the index.Condition— The used condition.Description— The index description (currently only used forSkipindexes).Parts— The number of parts after/before the index is applied.Granules— The number of granules after/before the index is applied.Ranges— The number of granules ranges after the index is applied.
Example:
"Node Type": "ReadFromMergeTree",
"Indexes": [
{
"Type": "Partition Min-Max",
"Keys": ["y"],
"Condition": "(y in [1, +inf))",
"Parts": 4/5,
"Granules": 11/12
},
{
"Type": "Partition",
"Keys": ["y", "bitAnd(z, 3)"],
"Condition": "and((bitAnd(z, 3) not in [1, 1]), and((y in [1, +inf)), (bitAnd(z, 3) not in [1, 1])))",
"Parts": 3/4,
"Granules": 10/11
},
{
"Type": "PrimaryKey",
"Keys": ["x", "y"],
"Condition": "and((x in [11, +inf)), (y in [1, +inf)))",
"Parts": 2/3,
"Granules": 6/10,
"Search Algorithm": "generic exclusion search"
},
{
"Type": "Skip",
"Name": "t_minmax",
"Description": "minmax GRANULARITY 2",
"Parts": 1/2,
"Granules": 2/6
},
{
"Type": "Skip",
"Name": "t_set",
"Description": "set GRANULARITY 2",
"": 1/1,
"Granules": 1/2
}
]With projections = 1, the Projections key is added. It contains an array of analyzed projections. Each projection is described as JSON with following keys:
Name— The projection name.Condition— The used projection primary key condition.Description— The description of how the projection is used (e.g. part-level filtering).Selected Parts— Number of parts selected by the projection.Selected Marks— Number of marks selected.Selected Ranges— Number of ranges selected.Selected Rows— Number of rows selected.Filtered Parts— Number of parts skipped due to part-level filtering.
Example:
"Node Type": "ReadFromMergeTree",
"Projections": [
{
"Name": "region_proj",
"Description": "Projection has been analyzed and is used for part-level filtering",
"Condition": "(region in ['us_west', 'us_west'])",
"Search Algorithm": "binary search",
"Selected Parts": 3,
"Selected Marks": 3,
"Selected Ranges": 3,
"Selected Rows": 3,
"Filtered Parts": 2
},
{
"Name": "user_id_proj",
"Description": "Projection has been analyzed and is used for part-level filtering",
"Condition": "(user_id in [107, 107])",
"Search Algorithm": "binary search",
"Selected Parts": 1,
"Selected Marks": 1,
"Selected Ranges": 1,
"Selected Rows": 1,
"Filtered Parts": 2
}
]With actions = 1, added keys depend on step type.
Example:
EXPLAIN json = 1, actions = 1, description = 0 SELECT 1 FORMAT TSVRaw;[
{
"Plan": {
"Node Type": "Expression",
"Node Id": "Expression_5",
"Expression": {
"Inputs": [
{
"Name": "dummy",
"Type": "UInt8"
}
],
"Actions": [
{
"Node Type": "INPUT",
"Result Type": "UInt8",
"Result Name": "dummy",
"Arguments": [0],
"Removed Arguments": [0],
"Result": 0
},
{
"Node Type": "COLUMN",
"Result Type": "UInt8",
"Result Name": "1",
"Column": "Const(UInt8)",
"Arguments": [],
"Removed Arguments": [],
"Result": 1
}
],
"Outputs": [
{
"Name": "1",
"Type": "UInt8"
}
],
"Positions": [1]
},
"Plans": [
{
"Node Type": "ReadFromStorage",
"Node Id": "ReadFromStorage_0"
}
]
}
}
]With compact = 0 and actions = 1, the Expression steps can be seen along with detailed information about expressions:
EXPLAIN actions = 1, compact = 0 SELECT sum(number) FROM numbers(10) GROUP BY number % 4;Output: sum(number)
Expression ((Project names + Projection))
│ Actions: INPUT : 0 -> sum(__table1.number) UInt64 : 0
│ INPUT :: 1 -> modulo(__table1.number, 4_UInt8) UInt8 : 1
│ ALIAS sum(__table1.number) :: 0 -> sum(number) UInt64 : 2
│ Positions: 2
└──Aggregating
│ Keys: number MOD 4
│ Aggregates: sum(number)
│ Skip merging: 0
└──Expression ((Before GROUP BY + Change column names to column identifiers))
│ Actions: INPUT : 0 -> number UInt64 : 0
│ COLUMN Const(UInt8) -> 4_UInt8 UInt8 : 1
│ ALIAS number :: 0 -> __table1.number UInt64 : 2
│ FUNCTION modulo(__table1.number : 2, 4_UInt8 :: 1) -> modulo(__table1.number, 4_UInt8) UInt8 : 0
│ Positions: 0 2
└──ReadFromSystemNumbers
Output: numberWith distributed = 1, the output includes not only the local query plan but also the query plans that will be executed on remote nodes. This is useful for analyzing and debugging distributed queries.
Example with distributed table:
EXPLAIN distributed=1 SELECT * FROM remote('127.0.0.{1,2}', numbers(2)) WHERE number = 1;Union
Expression ((Project names + (Projection + (Change column names to column identifiers + (Project names + Projection)))))
Filter ((WHERE + Change column names to column identifiers))
ReadFromSystemNumbers
Expression ((Project names + (Projection + Change column names to column identifiers)))
ReadFromRemote (Read from remote replica)
Expression ((Project names + Projection))
Filter ((WHERE + Change column names to column identifiers))
ReadFromSystemNumbersExample with parallel replicas:
SET enable_parallel_replicas = 2, max_parallel_replicas = 2, cluster_for_parallel_replicas = 'default';
EXPLAIN distributed=1 SELECT sum(number) FROM test_table GROUP BY number % 4;Expression ((Project names + Projection))
MergingAggregated
Union
Aggregating
Expression ((Before GROUP BY + Change column names to column identifiers))
ReadFromMergeTree (default.test_table)
ReadFromRemoteParallelReplicas
BlocksMarshalling
Aggregating
Expression ((Before GROUP BY + Change column names to column identifiers))
ReadFromMergeTree (default.test_table)In both examples, the query plan shows the complete execution flow including local and remote steps.
With pretty = 1, the plan tree is displayed using line-drawing characters instead of indentation, and additional information is shown for key steps:
- Query output columns are printed at the top of the plan.
- Expressions in filters, aggregation keys, sort descriptions, and window functions are displayed in human-readable SQL-like notation (e.g.,
a + 1 > 5instead ofgreater(plus(a, 1), 5)). Internal column identifier prefixes (such as__table1.) are removed for clarity. - Source steps (such as
ReadFromMergeTree) display their output columns. - Filter steps display the filter condition in SQL notation. When runtime join filters are present, they are shown separately.
- Aggregation steps display keys and aggregate functions with their arguments (e.g.,
sum(c),count()). - IN sets from tuple literals show their values (truncated for large sets), subquery-based sets are labeled
subquery1,subquery2, etc., and sets fromSetengine tables show the table name. - Join steps display the join relation using mathematical notation, the estimates the join-order optimizer produced for the step (cost, selectivity, output rows, and per-side rows), and the input columns of each side. The following symbols are used to represent different join types:
| Symbol | Join Type |
|---|---|
⋈ |
Inner Join |
⟕ |
Left Join |
⟖ |
Right Join |
⟗ |
Full Join |
⋉ |
Left Semi Join |
⋊ |
Right Semi Join |
⋉ with strikethrough |
Left Anti Join |
⋊ with strikethrough |
Right Anti Join |
× |
Cross Join |
For example, t1 ⟕ t2 means a left join between tables t1 and t2.
The number in brackets after the table name (e.g., t1[100]) indicates the estimated row count
when table statistics are available.
Below the join relation, each join step prints the estimates the join-order optimizer produced for it:
Cost: estimated <cost>
Selectivity: estimated (NDV) <selectivity>
Output rows: estimated <rows>
Left: rows estimated <left_rows>
Right: rows estimated <right_rows>Cost— the cost of the whole join subtree under this step, which is the value the optimizer minimizes when it compares candidate join orders. The cost of a join is its estimated number of matched row pairs,<selectivity> * <left_rows> * <right_rows>, plus the cost of its inputs.Selectivity— the estimated fraction of the Cartesian product of the two sides that survives the join condition. It is derived from the number of distinct values (NDV) of the join keys: a key equality keeps about1 / max(NDV_left, NDV_right)of the pairs, and the smallest fraction over the join conditions is used.Output rows— the estimated number of rows the join produces:<selectivity> * <left_rows> * <right_rows>for an inner join, floored at<left_rows>forLEFT, at<right_rows>forRIGHT, and at<left_rows> + <right_rows>forFULL, because an outer join keeps every row of its preserved side. When aSEMIorANTIjoin takes part in the reordering, it is estimated as a fraction of its preserved side,<preserved_rows> * min(1, <selectivity> * <other_rows>)forSEMIand the remaining rows forANTI; otherwise it uses the formula of its join kind.Left/Right— the estimated number of rows entering the join from each side.
A value the optimizer could not estimate is reported as no stats. This happens when the
join-order optimization did not run — for example, when
query_plan_optimize_join_order_limit
is 0 — or when there is no basis for the estimate.
The Input (left): and Input (right): lines list the columns each side feeds into the join.
The pretty option works well together with compact = 1, which hides Expression steps and detailed action info, making the plan easier to read.
A detailed example with joins. The
join_runtime_filter_min_probe_rows
setting is lowered only so that a table this small still builds a runtime join filter:
SET join_runtime_filter_min_probe_rows = 10;
CREATE TABLE t1 (id UInt64, value String) ENGINE = MergeTree ORDER BY id;
CREATE TABLE t2 (id UInt64, value String) ENGINE = MergeTree ORDER BY id;
INSERT INTO t1 SELECT number, toString(number) FROM numbers(100);
INSERT INTO t2 SELECT number, toString(number) FROM numbers(100);
EXPLAIN actions = 1, compact = 1, pretty = 1
SELECT * FROM t1 INNER JOIN t2 ON t1.id = t2.id FORMAT Raw;Output: id, value, id, value
Join (JOIN FillRightFirst)
│ t1[100] ⋈ t2[100]
│ Type: inner | Strictness: all | Algorithm: SpillingHashJoin(HashJoin)
│ Cost: estimated 100.00
│ Selectivity: estimated (NDV) 0.01
│ Output rows: estimated 100.00
│ Left: rows estimated 100.00
│ Right: rows estimated 100.00
│ Join conditions: id = id
│ Input (left): id, value
│ Input (right): id, value
├──ReadFromMergeTree (default.t1)
│ Read type: Default
│ Parts: 1 | Granules: 1
│ Output: id, value
│ Runtime filters: RF1(id, id from default.t2)
└──BuildRuntimeFilter (Build runtime join filter on id)
│ Filter id: RF1
│ Source table: default.t2
└──ReadFromMergeTree (default.t2)
Read type: Default
Parts: 1 | Granules: 1
Output: id, valueEXPLAIN PIPELINE
Settings:
header— Prints header for each output port. Default: 0.graph— Prints a graph described in the DOT graph description language. Default: 0.compact— Prints graph in compact mode ifgraphsetting is enabled. Default: 1.compact_repeated_processor_chains— Compacts adjacent repeated processor chains in text output by showing one copy of the chain with a repetition count. This can make parallel pipelines easier to read when the same chain appears many times, for example in joins. It does not affect graph output. Default: 0.
Resize 16 → 1
FillingRightJoinSide │
SimpleSquashingTransform │ × 16
Resize 1 → 16When compact=0 and graph=1 processor names will contain an additional suffix with unique processor identifier.
Example:
EXPLAIN PIPELINE SELECT sum(number) FROM numbers_mt(100000) GROUP BY number % 4;(Union)
(Expression)
ExpressionTransform
(Expression)
ExpressionTransform
(Aggregating)
Resize 2 → 1
AggregatingTransform × 2
(Expression)
ExpressionTransform × 2
(SettingQuotaAndLimits)
(ReadFromStorage)
NumbersRange × 2 0 → 1EXPLAIN ANALYZE
EXPLAIN ANALYZE actually runs the query, discards the result rows, and prints the same plan tree as EXPLAIN PLAN with each step annotated by what really happened at run time.
Settings:
EXPLAIN ANALYZE accepts the same display options as EXPLAIN PLAN (documented in the EXPLAIN PLAN section).
header— see EXPLAIN PLAN section.description— see EXPLAIN PLAN section.projections— see EXPLAIN PLAN section.sorting— see EXPLAIN PLAN section.input_headers— see EXPLAIN PLAN section.column_structure— see EXPLAIN PLAN section.actions— see EXPLAIN PLAN section. Default: 1.indexes— see EXPLAIN PLAN section. Default: 1.compact— see EXPLAIN PLAN section. Default: 1.pretty— see EXPLAIN PLAN section. Default: 1.processors— ForEXPLAIN ANALYZE, prints an additional line per stage with the per-processor elapsed time distribution:min,median,max, andsum. Useful to spot load skew across parallel processors. Default: 0.matches— ForEXPLAIN ANALYZE, makes join steps do the extra bookkeeping needed for thematched,match rateandfanoutmetrics in the cases where those numbers cannot be derived from what the join produces anyway. Where they can, they are reported without this option. See Join steps. Default: 0.
Example:
EXPLAIN ANALYZE SELECT number % 10 AS k, count() FROM numbers_mt(1000000) GROUP BY k;Query summary:
Time: 10.72 ms (planning 6.45 ms · execution 4.26 ms)
Read: 1.00 million rows, 8.00 MB (234.49 million rows/s., 1.88 GB/s.)
Peak memory: 28.98 KiB
Output: number MOD 10, count()
Expression ((Project names + Projection))
│ I/O: rows 10 → 10 · 90 B → 90 B
│ time 21.82 us (0.5%) · parallelism 0.98/1
└──Aggregating
│ Keys: number MOD 10
│ Aggregates: count()
│ Skip merging: 0
│ I/O: rows 1.00 million → 10 (0.00%) · 1.00 MB → 90 B
│ Stage (partial aggregation): time 868.45 us (20.4%) · parallelism 3.80/15
│ Stage (final aggregation): time 445.27 us (10.4%) · parallelism 1.11/16
└──Expression ((Before GROUP BY + Change column names to column identifiers))
│ I/O: rows 1.00 million → 1.00 million · 8.00 MB → 1.00 MB
│ time 677.07 us (15.9%) · parallelism 4.31/15
└──ReadFromSystemNumbers
Output: number
I/O: rows 0 → 1.00 million · 0 B → 8.00 MB
time 993.94 us (23.3%) · parallelism 7.52/15Let’s examine the output. First let’s look at the header.
Query summary:
Time: <total> (planning <planning> · execution <execution>)
Read: <rows> rows, <bytes> (<rows/s>, <bytes/s>)
Peak memory: <peak>Time— total time split into planning (i.e. creation of plan + optimization of plan + pipeline construction) and execution (running the pipeline) phases.Read— rows and uncompressed bytes read from tables, with throughput - the same numbers the normal query footer reports as “Processed”.Peak memory— peak memory the query used.
Now let’s look at the new lines that appear in the query plan.
I/O: rows <in> → <out> (<selectivity>%) · <bytes_in> → <bytes_out>
[Stage (<stage>): ]time <t> (<share>%) · parallelism <avg>/<max>Rows and bytes are reported once for the whole step (the I/O line). Time and parallelism are reported per stage of the step on the following indented line(s).
rows <in> → <out>— rows that entered and left the step; (<selectivity>%) shows how much the step filtered (out/in) or expanded the data, it is hidden when input rows equals output rows and when input rows equals0.<bytes_in> → <bytes_out>— uncompressed in-memory bytes flowing through the step (omitted when both are zero).time <t> (<share>%)— wall-clock time the stage was active, and its share of query execution time (i.e. without build time). Note shares can add up to more than 100% because stages and steps run concurrently.parallelism <avg>/<max>— average number of CPU threads working within this stage at once, out of the maximum it could use. A value near max means the stage was well parallelized; near 1 means it ran mostly serially.Stage (<stage>)— the name of the stage. A step with a single stage prints the time line directly, without aStage (...)label. Steps with several stages print one labeled line per stage, e.g.AggregatingshowsStage (partial aggregation)andStage (final aggregation), and a hash join showsStage (build)andStage (probe).
Join steps
For a join step EXPLAIN ANALYZE prints lines comparing the join-order optimizer’s estimates with what actually happened (see Estimated vs. actual join metrics) and per-side participation lines — Left and Right — followed by any lines specific to the join implementation. Every value of join_algorithm is covered (hash, parallel_hash, grace_hash, partial_merge, full_sorting_merge, parallel_full_sorting_merge, direct), and so are the two implementations that setting cannot select: a CROSS or COMMA join and any ON section without a key equality, and the Join table engine. Most of them report both sides; some report only the side they materialize (for example direct prints only Left:).
The per-side lines share the same shape:
Left: rows estimated <estimated_left_rows> · rows <left_rows> · matched <matched_left_rows> · match rate <match_rate>% · fanout <fanout>
Right: rows estimated <estimated_right_rows> · rows <right_rows> · matched <matched_right_rows> · match rate <match_rate>% · fanout <fanout>For each side EXPLAIN ANALYZE reports:
rows estimated <estimated_rows>— the join-order optimizer’s estimate of that side’s rows, printed for comparison with the actualrowsnext to it;no statswhen the optimizer produced no estimate (see Estimated vs. actual join metrics).rows <rows>— the total number of rows of that side that passed through the join.matched <matched_rows>— the number of rows of that side that found at least one join partner on the other side. This counts rows, not keys: if a key occurs three times on the right and matches, all three right rows count as matched.match rate <match_rate>%— the percentage of that side’s rows that matched, computed as100 * <matched_rows> / <rows>.fanout <fanout>— how many output rows an average matched row of that side produced.
A number that cannot be derived exactly is reported as not collected rather than as 0. match rate and fanout are derived from matched, so a side without it reports all three as not collected.
Estimated vs. actual join metrics
A join step carries the join-order optimizer’s estimates — the same ones EXPLAIN PLAN shows (see the EXPLAIN PLAN section) — and EXPLAIN ANALYZE prints each of them next to the measured value:
Cost: estimated <cost> · actual <cost>
Selectivity: estimated (NDV) <selectivity> · actual (cartesian) <selectivity>
Output rows: estimated <rows> · actual <rows> · q-error <ratio>Cost— both values count matched output rows. The estimate is the optimizer’s cost of the join subtree:<selectivity> * <left_rows> * <right_rows>plus the cost of its inputs. The actual value is measured the same way — the matched output rows of this join plus the actual cost of every join below it that belongs to the same reorder cluster.Selectivity— the estimate is derived from the number of distinct values of the join keys; the actual value is the measured fraction of the Cartesian product that ended up in the output:<matched output rows> / (<left rows> * <right rows>), since this is what optimizer tries to estimate.Output rows— the estimated and the actual number of rows the join produced. When both are non-zero,q-errorreportsmax(estimated / actual, actual / estimated), the standard measure of cardinality-estimation quality:1.00means a perfect estimate, and a large value means the optimizer picked the join order using a badly wrong cardinality.
An estimate that was never made is reported as no stats — for example, when the join-order optimization did not run because query_plan_optimize_join_order_limit is 0. This is distinct from not collected, which marks an actual value the execution could not measure. Joins with a pre-filled right side (the Join table engine, direct joins) do not go through the join-order optimizer and print only the participation lines.
The join step also prints Input (left): and Input (right): lines with the columns each side feeds into the join.
For the tables of the EXPLAIN PLAN join example, the join step of EXPLAIN ANALYZE SELECT * FROM t1 INNER JOIN t2 ON t1.id = t2.id looks like this:
Join (JOIN FillRightFirst)
│ t1[100] ⋈ t2[100]
│ Type: inner | Strictness: all | Algorithm: SpillingHashJoin(HashJoin)
│ Join conditions: id = id
│ Cost: estimated 100.00 · actual 100.00
│ Selectivity: estimated (NDV) 0.01 · actual (cartesian) 0.01
│ Output rows: estimated 100.00 · actual 100.00 · q-error 1.00
│ Left: rows estimated 100.00 · rows 100.00 · matched 100.00 · match rate 100.00% · fanout 1.00
│ Right: rows estimated 100.00 · rows 100.00 · matched not collected · match rate not collected · fanout not collected
│ Hash table: unique keys 100.00 · memory 6.27 KB
│ Input (left): id, value
│ Input (right): id, value
│ I/O: rows 200 → 100 (50.00%) · 3.58 KB → 3.58 KB
│ Stage (build): time 127.45 us (1.2%) · parallelism 0.99/1
│ Stage (probe): time 124.99 us (1.2%) · parallelism 0.99/1Fanout
fanout measures row multiplication:
matched output rows = <output_rows> - <NULL-padded rows of both sides>
fanout = <matched output rows> / <matched_rows of that side>An outer join emits one NULL-padded output row for every row of a preserved side that found no partner. Those rows are subtracted so that they do not dilute the ratio. Only a preserved side has them — the right side for RIGHT and FULL, the left one for LEFT and FULL:
fanout = 0— the matched rows produced no output row at all, which is what anANTIjoin does: it emits only the rows that found no partner.fanout = 1— a clean 1:1 join; every matched row produced exactly one output row.fanout > 1— a 1:N join; duplicate keys on the other side multiplied the rows. A large value on both sides at once is the signature of an unintended Cartesian blowup.
When the numbers require matches = 1
Most of these numbers fall out of data the join builds anyway and are reported by a plain EXPLAIN ANALYZE. The rest need bookkeeping the join would otherwise not do, so they are only reported with EXPLAIN ANALYZE matches = 1. Which ones those are depends on the algorithm; in the hash family they are two cases:
- the right side of
ALL INNERandALL LEFT, which requires marking every matched right row; - the left side of
ALL LEFTandALL FULL, but only when the query selects nothing from the right table and theONsection is a plain key equality. Otherwise the probe already records which left rows matched — either to materialize the right columns or to evaluate the residual condition — and the count is exact without the option.
partial_merge needs it for the right side of the four ALL kinds, for the same reason.
full_sorting_merge and parallel_full_sorting_merge need it for both sides of the ANY kinds. The
ALL kinds need nothing.
matches = 1 does not make every combination collectable. Which side a join can report follows
from what that join has to do anyway, so it depends on the algorithm as well as on the kind and
strictness.
Hash family. hash, parallel_hash and grace_hash always agree with each other:
| Join | matched left |
matched right |
|---|---|---|
ALL INNER, ALL LEFT, ALL RIGHT, ALL FULL |
yes | yes |
SEMI LEFT, ANTI LEFT |
yes | no |
ANY RIGHT, ANTI RIGHT |
no | yes |
ASOF (inner) |
yes | no |
SEMI RIGHT |
no | no |
ANY INNER, ANY LEFT, ASOF LEFT |
no | no |
The right side is unavailable whenever the join keeps only one row per key in its hash table, which ANY, SEMI and ANTI joins do: the duplicate right rows are never stored, so they cannot be counted. The left side is unavailable when the join suppresses the output of a left row whose partner was already claimed by another left row, which makes the emitted rows an undercount of the matched ones.
Enabling any_join_distinct_right_table_keys switches ANY to the older RightAny semantics, which emits one row per left row and therefore keeps both counts. ANY RIGHT and ANY FULL then report both sides, and ANY INNER is rewritten to SEMI LEFT.
The Join table engine follows the same table, using the kind and strictness declared in the engine: Join(ALL, INNER, …) reports both sides, Join(ANY, LEFT, …) neither.
Merge algorithms. full_sorting_merge and parallel_full_sorting_merge accept the four ALL kinds, ANY INNER, ANY LEFT, ANY RIGHT, ASOF and ASOF LEFT. They report both sides for every kind except ASOF and ASOF LEFT, where the right side is not collected, and without matches = 1 — they walk the two sorted inputs and see every row of an equal range as they consume it, so nothing has to be reconstructed afterwards.
partial_merge accepts ALL INNER, ALL LEFT, ALL RIGHT, ALL FULL, ANY INNER, ANY LEFT and SEMI LEFT. It reports both sides for the four ALL kinds, the right one with matches = 1; for ANY INNER, ANY LEFT and SEMI LEFT the right side is not collected.
direct. The left side only. The right side is a key-value store that is never materialized into rows, so it has no Right: line at all.
CROSS, COMMA and a constant ON. Neither side, as described above.
Where two algorithms both report a number, the numbers agree. The merge algorithms simply have more information; they do not disagree about what a match is.
Algorithm-specific lines
Let’s take a look at the lines each join implementation adds on top of those.
For hash and parallel_hash joins, and for the Join table engine, a Hash table: line describes the hash table built from the right table:
Hash table: unique keys <unique_keys> · memory <peak_memory>unique keys <unique_keys>— the number of unique keys stored in the hash table during the build phase.memory <peak_memory>— the peak memory used by the hash table during the build phase.
For grace_hash join the Hash table: line additionally reports how the join adapted to the memory limit, and a Spill: line reports whether data was spilled to disk:
Hash table: unique keys <unique_keys> · memory <peak_memory> · buckets <buckets> · rehashes <rehashes>
Spill: yes · left spilled <left_spilled_bytes> · right spilled <right_spilled_bytes>buckets <buckets>— the number of buckets the grace hash join ended up with by the end of execution. This is always a power of2.rehashes <rehashes>— how many times the number of buckets had to be doubled in order to fit into the memory limit.Spill:— ayes/noflag telling whether any spilling to disk happened. When it did,left spilled <left_spilled_bytes>andright spilled <right_spilled_bytes>report the compressed bytes spilled from the left (probe) and right (build) sides; when nothing was spilled the line is simplySpill: no.
For partial_merge join the Right: line carries extra information about how the right table was buffered and sorted, and the sorting time is shown on the Stage (build) and Stage (probe) lines:
Right: rows estimated <estimated_right_rows> · rows <right_rows> · matched <matched_right_rows> · size <right_size> · blocks <right_blocks> · storage <in-memory|external> · match rate <match_rate>% · fanout <fanout>
Stage (build): time <t> (<share>%) · parallelism <avg>/<max> · sort time <build_sort_time> · sort share <build_sort_share>%
Stage (probe): time <t> (<share>%) · parallelism <avg>/<max> · sort time <probe_sort_time> · sort share <probe_sort_share>%size <right_size>— the memory size of the blocks of right table.blocks <right_blocks>— the number of blocks the right table was buffered into.storage <in-memory|external>— whether the right table fit into memory (in-memory) or had to be spilled to disk (external). When it isexternal, an extraspilled <spilled_bytes>reports the compressed bytes written to disk.sort time <sort_time>— the time spent sorting the right table (on the build stage) and each incoming left block (on the probe stage).sort share <sort_share>%—sort timeas a share of that stage’s own busy time (the sum of its processors’ elapsed time), unlike the stagetimepercentage, which is a share of the whole query’s execution time.
For full_sorting_merge join only the common Left: and Right: lines are printed.
For direct join only the Left: line is printed, since the right side is a key-value store that is looked up directly rather than materialized into rows.
For a CROSS or COMMA join, and for any ON section without a key equality, a Buffer: line describes how the right table was held in memory and a Spill: line reports whether it went to disk:
Buffer: memory <peak_memory> · compressed <yes|no>
Spill: yes · right spilled <right_spilled_bytes>memory <peak_memory>— the peak memory the buffered right table occupied.compressed <yes|no>— whether at least one buffered block was compressed; readers then decompress every stored block.Spill:— the sameyes/noflag as forgrace_hash, withright spilled <right_spilled_bytes>reporting the compressed bytes written to disk.
Both sides report matched not collected here: a constant predicate either pairs every left row with every right row or with none, so asking which individual rows matched has no answer.
For a join against the Join table engine both sides are reported, together with the Hash table: line describing the pre-built table. The right side counts the rows stored in the engine, not the rows of some per-query build.
Per-processor times
With processors = 1, an extra line is printed under each stage, showing the distribution of elapsed time across the stage’s processors:
Time per processor (<n>): min <t> · median <t> · max <t> · sum <t><n> is the number of processors in the stage. A large gap between median and max points to load skew between parallel processors.
EXPLAIN ESTIMATE
Shows the estimated number of rows, marks and parts to be read from the tables while processing the query. Works with tables in the MergeTree family.
Example
Creating a table:
CREATE TABLE ttt (i Int64) ENGINE = MergeTree() ORDER BY i SETTINGS index_granularity = 16, write_final_mark = 0;
INSERT INTO ttt SELECT number FROM numbers(128);
OPTIMIZE TABLE ttt;EXPLAIN ESTIMATE SELECT * FROM ttt;┌─database─┬─table─┬─parts─┬─rows─┬─marks─┐
│ default │ ttt │ 1 │ 128 │ 8 │
└──────────┴───────┴───────┴──────┴───────┘EXPLAIN WHATIF
Estimates the benefit a hypothetical skip index would have on a SELECT query, without materializing the index on disk. Define one or more candidates with CREATE HYPOTHETICAL INDEX, then run EXPLAIN WHATIF SELECT ... to see, for each candidate: applicability, estimated marks read, estimated bytes, and skip ratio.
Hypothetical projections defined with CREATE HYPOTHETICAL PROJECTION are candidates too. A normal projection is estimated by building its primary index in memory over the parts the query would read and pruning it as a materialized projection would be pruned. The report gives the marks and rows the projection read would touch, a read_ratio against the base-table read (below 1x means less work, above means more) and a verdict with the reason behind it, following the optimizer’s rule: the projection wins when it reads fewer marks than the base table, or the same number while serving an outer ORDER BY. Listed as status: not_applicable and not estimated yet: aggregate projections, projections with a WHERE clause or their own skip indexes (WITH SETTINGS add_minmax_index_*), projections ordered by the commit order (TYPE commit_order and its _block_number, _block_offset query form), projections whose sort key is not among the columns they store (for example ORDER BY _part_offset), projections that store _block_number (the writer builds those only when a part is merged), and projections that do not provide every column the query reads. force_optimize_projection, force_optimize_projection_name and preferred_optimize_projection_name are ignored. The mark count is modelled by sizing granules the way the writer does, one granule size per block the writer is handed. Which blocks that is depends on the path that writes the projection part - one squashed block for an insert or a materialization, runs of merge_max_block_size for a merge - and a part records none of it, so the estimate is computed for each of those layouts. When they do not agree on the comparison with the base read, the report gives the range and verdict: too close to call instead of a decision. A projection whose definition no longer fits the table is reported with that reason.
Syntax
EXPLAIN WHATIF [empirical = 0] SELECT ...Settings
empirical—1(default) runs the index over the baseline-pruned granules in memory to measure the skip ratio (an upper bound).0skips that path. Either way, if empirical doesn’t produce a result (disabled, or the index can’t be evaluated in memory) the estimator falls back to column statistics, and finally to an applicability-only summary if neither is available.
Output
Baseline (after PK + partition + existing indexes):
table: db.t
parts: 1
marks: 100
rows: 10000
est_bytes: 1.50 MiB (only when the query reads rows)
With idx_b (minmax, hypothetical):
status: applicable
marks: 1
est_bytes: 15.00 KiB (only when baseline bytes are known)
skip_ratio: 99.0%
Estimation:
source: empirical | statistical | applicability_only
empirical_status: ok | unsupported | disabled
empirical_reason: <reason> (only when empirical_status = unsupported)
sampled_parts: 50 / 100 (only when source = empirical)
sampled_marks: 50 / 100 (only when source = empirical)
elapsed_us: 631 (only when source = empirical)source— how the estimate was produced.empirical: built the index in memory over the baseline-pruned granules and counted the granules the index would skip. This is an upper bound — see the limitations inCREATE HYPOTHETICAL INDEX.statistical: derived from column statistics. Used when empirical is disabled (empirical = 0) or empirical couldn’t produce a result, and column statistics are defined on the relevant columns.applicability_only: the index is applicable to the predicate but neither empirical nor statistical estimation produced a result (e.g.empirical = 0and no column statistics defined). Reportsskip_ratio: 0.0%as a conservative bound.
empirical_reason— why the empirical estimate could not run. Shown only withempirical_status: unsupported. For example, a non-zeromerge_tree_min_rows_for_seekormerge_tree_min_bytes_for_seekmakes a real read coalesce mark ranges, which the per-granule count does not model, so the estimate falls back tostatisticalorapplicability_only.sampled_parts/sampled_marks—<baseline-pruned> / <total in the table>. Shows what fraction of the table survived PK, partition, and existing-index pruning, i.e. the input to the hypothetical index.est_bytes— an estimate of the bytes read, derived from the table’s average row size, so it is approximate and varies with storage and compression. The baseline line appears only when the query reads rows; the per-candidate line only when the baseline byte estimate is known.
The setting is written inline between WHATIF and the SELECT — there is no SETTINGS keyword (this matches how other EXPLAIN variants accept their options).
If neither hypothetical indexes nor hypothetical projections are defined for the table, EXPLAIN WHATIF reports status: not_applicable with a hint to create one.
Combined row (multiple candidates)
When two or more candidates are evaluated empirically, EXPLAIN WHATIF appends one extra block named (combined: idx_a, idx_b, ...) after the per-candidate rows. It reports the joint benefit of having all of those indexes at once: a real read keeps a granule only if it survives every skip index, so the combined estimate is the intersection of the candidates’ surviving granules. Its skip_ratio is therefore at least as high as the best single candidate — complementary indexes prune more together, while redundant ones leave it unchanged.
Only candidates with source: empirical contribute, because the combined row is built by intersecting their per-granule survival sets. Candidates estimated statistical or applicability_only have no per-granule data and are excluded; consequently the combined block appears only when at least two candidates produced an empirical estimate, and is omitted otherwise (for example under empirical = 0). Its estimation fields read the same as a per-candidate empirical block, except elapsed_us is 0 — the combined estimate is derived from the per-candidate scans, not a new scan. The synthetic (combined: ...) name is a report label only and cannot be used with force_data_skipping_indices.
Empirical example
CREATE TABLE t (a UInt64, b UInt64) ENGINE = MergeTree ORDER BY a
SETTINGS index_granularity = 100;
INSERT INTO t SELECT number, number FROM numbers(10000);
CREATE HYPOTHETICAL INDEX idx_b ON t (b) TYPE minmax GRANULARITY 1;
EXPLAIN WHATIF SELECT * FROM t WHERE b = 42;Baseline (after PK + partition + existing indexes):
table: default.t
parts: 1
marks: 100
rows: 10000
est_bytes: 85.52 KiB
With idx_b (minmax, hypothetical):
status: applicable
marks: 1
est_bytes: 875.00 B
skip_ratio: 99.0%
Estimation:
source: empirical
empirical_status: ok
sampled_parts: 1 / 1
sampled_marks: 100 / 100The hypothetical minmax would prune from 100 marks down to 1 — skip_ratio: 99.0%. (est_bytes is an estimate from the average row size, so the exact figure varies.)
Statistical example
Column statistics are off by default. To exercise the statistical path, define them on the relevant columns first and wait for the materialize mutation to finish:
ALTER TABLE t ADD STATISTICS b TYPE tdigest;
ALTER TABLE t MATERIALIZE STATISTICS b SETTINGS mutations_sync = 1;Then disable the empirical path so the estimator falls back to column statistics:
EXPLAIN WHATIF empirical = 0 SELECT * FROM t WHERE b < 10;With idx_b (minmax, hypothetical):
status: applicable
marks: 1
est_bytes: 1.66 KiB
skip_ratio: 99.9%
Estimation:
source: statistical
empirical_status: disabledThe number comes from the column-statistic selectivity of b < 10 (about 10 rows out of 10000) and is reported as an upper bound on skip_ratio. There are no sampled_parts / sampled_marks — no data was read.
If neither path is available (e.g. empirical = 0 and no column statistics defined), the estimator reports source: applicability_only and a conservative skip_ratio: 0.0%.
Hypothetical projections
A normal projection is estimated the same way, against the parts the query would read after primary-key and partition pruning:
CREATE TABLE t (a UInt64, b UInt64, v UInt64) ENGINE = MergeTree ORDER BY a SETTINGS index_granularity = 100;
INSERT INTO t SELECT number, number % 100, number FROM numbers(10000);
CREATE HYPOTHETICAL PROJECTION p_b ON t (SELECT a, b, v ORDER BY b);
EXPLAIN WHATIF SELECT count() FROM t WHERE b = 42;Baseline (after PK + partition + existing indexes):
table: default.t
parts: 1
marks: 100
rows: 10000
est_bytes: 80.79 KiB
With p_b (normal projection, hypothetical):
status: applicable
marks: 2
rows: 200
read_ratio: 0.02x
verdict: chosen
reason: 2 marks would be read instead of 100 from the base table
Estimation:
source: empirical
empirical_status: ok
sampled_parts: 1 / 1
sampled_marks: 100 / 100
elapsed_us: 1126marks and rows are what the projection read itself would touch, not a share of the base-table read, because a projection granule holds different rows than a base granule. read_ratio is those marks over the base-table marks, so 0.02x is fifty times less work and 15x is fifteen times more. verdict applies the optimizer’s own rule — fewer marks than the base table, or the same number when the projection serves the query’s ORDER BY — so a candidate can be applicable and still not be chosen.
EXPLAIN TABLE OVERRIDE
Shows the result of a table override on a table schema accessed through a table function. Also does some validation, throwing an exception if the override would have caused some kind of failure.
Example
Assume you have a remote MySQL table like this:
CREATE TABLE db.tbl (
id INT PRIMARY KEY,
created DATETIME DEFAULT now()
)EXPLAIN TABLE OVERRIDE mysql('127.0.0.1:3306', 'db', 'tbl', 'root', 'clickhouse')
PARTITION BY toYYYYMM(assumeNotNull(created))┌─explain─────────────────────────────────────────────────┐
│ PARTITION BY uses columns: `created` Nullable(DateTime) │
└─────────────────────────────────────────────────────────┘