解读 EXPLAIN 计划
阅读 Explain 计划
简介
本节介绍如何阅读 DataFusion 的查询计划。虽然要完全理解这些计划的所有细节需要对 DataFusion 引擎有相当深入的专业知识,但本指南将帮助你入门基础知识。
DataFusion 使用 query plan(查询计划)来执行查询。若要在不实际运行查询的情况下查看该计划,可以在 SQL 查询中加上关键字 EXPLAIN,或调用 DataFrame::explain 方法。
示例:选择与过滤
在本节中,我们将针对 hits.parquet 文件运行示例查询。关于如何获取该文件,请参见下文。
我们来看 DataFusion 如何运行一个查询,选出站点 http://domcheloveplanet.ru/ 上排名前 5 的观看列表:
EXPLAIN FORMAT INDENT SELECT "WatchID" AS wid, "hits.parquet"."ClientIP" AS ip
FROM 'hits.parquet'
WHERE starts_with("URL", 'http://domcheloveplanet.ru/')
ORDER BY wid ASC, ip DESC
LIMIT 5;输出结果如下:
+---------------+-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
| plan_type | plan |
+---------------+-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
| logical_plan | Sort: wid ASC NULLS LAST, ip DESC NULLS FIRST, fetch=5 |
| | Projection: hits.parquet.WatchID AS wid, hits.parquet.ClientIP AS ip |
| | Filter: starts_with(hits.parquet.URL, Utf8("http://domcheloveplanet.ru/")) |
| | TableScan: hits.parquet projection=[WatchID, ClientIP, URL], partial_filters=[starts_with(hits.parquet.URL, Utf8("http://domcheloveplanet.ru/"))] |
| physical_plan | SortPreservingMergeExec: [wid@0 ASC NULLS LAST,ip@1 DESC], fetch=5 |
| | SortExec: TopK(fetch=5), expr=[wid@0 ASC NULLS LAST,ip@1 DESC], preserve_partitioning=[true] |
| | ProjectionExec: expr=[WatchID@0 as wid, ClientIP@1 as ip] |
| | CoalesceBatchesExec: target_batch_size=8192 |
| | FilterExec: starts_with(URL@2, http://domcheloveplanet.ru/) |
| | DataSourceExec: file_groups={16 groups: [[hits.parquet:0..923748528], ...]}, projection=[WatchID, ClientIP, URL], predicate=starts_with(URL@13, http://domcheloveplanet.ru/), file_type=parquet |
+---------------+------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
2 row(s) fetched.
Elapsed 0.060 seconds.它包含两个部分:逻辑计划与物理计划
- 逻辑计划(Logical Plan): 是针对特定的 SQL 查询、DataFrame 或其他语言生成的计划,此时并不了解底层的数据组织方式。
- 物理计划(Physical Plan): 是在逻辑计划的基础上,结合硬件配置(例如 CPU 数量)和底层数据组织(例如文件数量)生成的计划。该物理计划与你的硬件配置和数据是对应的。如果你将相同的数据加载到配置不同的其他硬件上,同一查询可能会生成不同的查询计划。
理解查询计划有助于你了解其性能。例如,当计划显示你的查询读取了大量文件时,这提示你要么在查询中增加更多过滤条件以读取更少的数据,要么调整文件设计,使文件数量更少但单个文件更大。本文档重点介绍如何阅读查询计划。至于如何让查询运行得更快,则取决于查询变慢的原因,超出了本文档的范围。
查询计划是一棵树
查询计划是一棵倒置的树,我们总是自底向上阅读。图 1 中的物理计划以树的形式表示如下:
▲
│
│
┌─────────────────────────────────────────────────┐
│ SortPreservingMergeExec │
│ [wid@0 ASC NULLS LAST,ip@1 DESC] │
│ fetch=5 │
└─────────────────────────────────────────────────┘
▲
│
┌─────────────────────────────────────────────────┐
│ SortExec TopK(fetch=5), │
│ expr=[wid@0 ASC NULLS LAST,ip@1 DESC], │
│ preserve_partitioning=[true] │
└─────────────────────────────────────────────────┘
▲
│
┌─────────────────────────────────────────────────┐
│ ProjectionExec │
│ expr=[WatchID@0 as wid, ClientIP@1 as ip] │
└─────────────────────────────────────────────────┘
▲
│
┌─────────────────────────────────────────────────┐
│ CoalesceBatchesExec │
└─────────────────────────────────────────────────┘
▲
│
┌─────────────────────────────────────────────────┐
│ FilterExec │
│ starts_with(URL@2, http://domcheloveplanet.ru/) │
└─────────────────────────────────────────────────┘
▲
│
┌────────────────────────────────────────────────┐
│ DataSourceExec │
│ hits.parquet (filter = ...) │
└────────────────────────────────────────────────┘树/计划中的每个节点都以 Exec 结尾,有时也被称为 operator 或 ExecutionPlan,数据在其中被处理、转换并向上层传递。
- 首先,
hits.parquet文件中的 Parquet 数据由DataSourceExec使用 16 个内核并行读取为 16 个「分区」(稍后会详细介绍),该节点在扫描过程中会进行第一轮过滤。 - 接下来,输出结果经过
FilterExec过滤,确保只有starts_with(URL, 'http://domcheloveplanet.ru/')求值为 true 的行才会继续向上传递。 - 然后
CoalesceBatchesExec确保数据被合并成更大的批次以便处理。 - 接着
ProjectionExec对数据进行投影,将WatchID和ClientIP列分别重命名为wid和ip。 - 然后
SortExec按照wid ASC, ip DESC对数据排序。Topk(fetch=5)表示这里使用了一种特殊实现,只在每个分区中追踪并输出前 5 个值。 - 最后,
SortPreservingMergeExec合并来自所有分区的已排序数据,并返回整体的前 5 行。
理解大型查询计划
一个大型查询计划可能看起来令人生畏,但只要遵循以下步骤,你就能快速理解它的作用:
- 和往常一样,自底向上逐个算子阅读。
- 通过阅读物理计划文档来理解该算子的职责。
- 理解该算子的输入数据,以及数据量可能有多大或多小。
- 理解该算子会产出多少数据,产出的数据是什么样的。
如果你能回答这些问题,就能估算出这个计划需要完成多少工作,进而知道它需要多长时间。不过,EXPLAIN 只是向你展示计划,并不会执行它。
如果你想知道查询计划中每个算子具体完成了多少工作,可以使用 EXPLAIN ANALYZE 获取带有运行时信息的执行计划说明(见下一节)。
更多调试信息:EXPLAIN VERBOSE
如果计划需要读取的文件太多,EXPLAIN 不会显示全部文件。要查看它们,请使用 EXPLAIN VERBOSE。与 EXPLAIN 一样,EXPLAIN VERBOSE 不会运行查询。相反,它会显示完整的执行计划说明,包括默认说明中省略的信息,以及 DataFusion 在返回结果之前生成的所有中间物理计划。这一模式对调试非常有用,可以查看 DataFusion 是何时以及为何向计划中添加或移除算子的。
执行计数器:EXPLAIN ANALYZE
在执行过程中,DataFusion 的算子会收集详细的指标信息。你可以通过 ExecutionPlan::metrics 以编程方式访问这些指标,也可以使用 EXPLAIN ANALYZE 命令。例如,下面是对上面同一条查询使用 EXPLAIN ANALYZE 的结果(注意:为清晰起见,输出经过了编辑)。
> EXPLAIN ANALYZE SELECT "WatchID" AS wid, "hits.parquet"."ClientIP" AS ip
FROM 'hits.parquet'
WHERE starts_with("URL", 'http://domcheloveplanet.ru/')
ORDER BY wid ASC, ip DESC
LIMIT 5;
+-------------------+----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
| plan_type | plan |
+-------------------+----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
| Plan with Metrics | SortPreservingMergeExec: [wid@0 ASC NULLS LAST,ip@1 DESC], fetch=5, metrics=[output_rows=5, elapsed_compute=2.375µs] |
| | SortExec: TopK(fetch=5), expr=[wid@0 ASC NULLS LAST,ip@1 DESC], preserve_partitioning=[true], metrics=[output_rows=75, elapsed_compute=7.243038ms, row_replacements=482] |
| | ProjectionExec: expr=[WatchID@0 as wid, ClientIP@1 as ip], metrics=[output_rows=811821, elapsed_compute=66.25µs] |
| | FilterExec: starts_with(URL@2, http://domcheloveplanet.ru/), metrics=[output_rows=811821, elapsed_compute=1.36923816s] |
| | DataSourceExec: file_groups={16 groups: [[hits.parquet:0..923748528], ...]}, projection=[WatchID, ClientIP, URL], predicate=starts_with(URL@13, http://domcheloveplanet.ru/), metrics=[output_rows=99997497, elapsed_compute=16ns, ... bytes_scanned=3703192723, ... time_elapsed_opening=308.203002ms, time_elapsed_scanning_total=8.350342183s, ...] |
+-------------------+----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
1 row(s) fetched.
Elapsed 0.720 seconds.在这种情况下,DataFusion 实际上执行了查询,但丢弃了所有结果,转而返回一个带有新增字段 metrics=[...] 的经过标注的计划。
大多数算子都有 output_rows 和 elapsed_compute 这两个通用指标,有些算子还具有特定于自身的指标,例如使用 ParquetSource 的 DataSourceExec 就有 bytes_scanned=3703192723。注意,时间与计数器是在所有核心上汇总统计的,因此如果你有 16 个核心,报告的时间就是这 16 个核心所耗时间的总和。
同样,从下往上阅读:
DataSourceExecoutput_rows=99997497:共产生了 9999 万行bytes_scanned=3703192723:在 14GB 的文件中,实际读取了 3.7GB(由于下推了投影)bytes_processed=14779976446:全部 14GB 都被计入:读取的 3.7GB,加上被裁剪掉的行组所占用的字节数。将其与计划中文件的总大小进行比较,就可以了解扫描的进度time_elapsed_opening=308.203002ms:打开文件并准备读取花费了 300 毫秒time_elapsed_scanning_total=8.350342183s:实际解码 Parquet 数据花费了 8.3 秒的 CPU 时间(跨 16 个核心)
FilterExecoutput_rows=811821:在其输入的 9999 万行中,只有 81.1 万行通过了过滤并输出elapsed_compute=1.36923816s:计算过滤条件总计花费了 1.36 秒的 CPU 时间(跨 16 个核心)
CoalesceBatchesExecoutput_rows=811821、elapsed_compute=12.873379ms:在 13 毫秒内产出了 81.1 万行
ProjectionExecoutput_rows=811821, elapsed_compute=66.25µs:在 66 微秒内产出了 81.1 万行。由于该投影不操作任何数据,因此几乎是瞬间完成的
SortExecoutput_rows=75:总计产出了 75 行。16 个核心中的每一个最多可以产出 5 行,但在本例中并非所有核心都产出了。elapsed_compute=7.243038ms:用于确定前 5 行的时间为 7 毫秒row_replacements=482:在内部,TopK 算子更新其顶部列表 482 次
SortPreservingMergeExecoutput_rows=5、elapsed_compute=2.375µs:在 2.375 微秒内产出了最终的 5 行
当启用谓词下推时,使用 ParquetSource 的 DataSourceExec 会获得以下指标:
output_rows_skew:由各分区的output_rows推导出的输出倾斜度得分。0%表示完全均衡,100%表示倾斜到最大程度,N/A表示没有产生任何输出行。page_index_rows_pruned:由页索引过滤器评估的行数。该指标同时报告总共考虑了多少行,以及有多少行匹配(未被裁剪)。page_index_pages_pruned:由页索引过滤器评估的页数。该指标同时报告总共考虑了多少页,以及有多少页匹配(未被裁剪)。page_index_pages_skipped_by_fully_matched:跳过页索引裁剪的页数,原因是行组统计信息已证明所属行组被完全匹配。这些页仍会被扫描。row_groups_pruned_bloom_filter:由布隆过滤器评估的行组数,报告检查的行组总数和匹配的行组数。row_groups_pruned_statistics:由行组统计信息(最小值/最大值)评估的行组数,报告检查的行组总数和匹配的行组数。limit_pruned_row_groups:被 limit 裁剪的行组数。limit_pruned_rows:当完全匹配的页范围包含足够的行以满足 limit 时,被 limit 裁剪跳过的行数。pushdown_rows_matched:经过上述任一过滤器测试并全部通过的行数。pushdown_rows_pruned:经过上述任一过滤器测试,但至少未通过其中一项的行数。predicate_evaluation_errors:过滤表达式求值失败的次数(正常运行时预期为零)num_predicate_creation_errors:创建谓词出错的次数(正常运行时预期为零)bloom_filter_eval_time:解析和求值布隆过滤器所花费的时间statistics_eval_time:解析和求值行组级统计信息所花费的时间row_pushdown_eval_time:求值行级过滤器所花费的时间page_index_eval_time:求值页索引过滤器所需的时间
Postgres 风格的 EXPLAIN (...) 选项
除了传统的关键字形式(EXPLAIN ANALYZE VERBOSE FORMAT tree SELECT ...)之外,对于 supports_explain_with_utility_options 返回 true 的方言,DataFusion 也接受 Postgres 风格的选项列表。这包括默认的 GenericDialect、PostgreSqlDialect 和 DuckDbDialect 等。
EXPLAIN (ANALYZE, VERBOSE, METRICS 'rows,bytes', LEVEL dev)
SELECT ... ;可识别的选项如下:
| 选项 | 参数 | 作用 |
|---|---|---|
ANALYZE | 布尔值,可选 | 执行计划并收集度量指标。单独书写时默认为 TRUE。等价于 ANALYZE 关键字。 |
VERBOSE | 布尔值,可选 | 显示每个分区的度量指标以及更多细节。等价于 VERBOSE 关键字。 |
FORMAT | 标识符/字符串 | indent、tree、pgjson、graphviz 之一。等价于 FORMAT <format> 子句。 |
METRICS | 字符串 | 按类别过滤 ANALYZE 的度量指标。接受 'all'、'none',或 rows,bytes,timing,uncategorized 的任意以逗号分隔的子集。 |
LEVEL | 标识符/字符串 | summary 或 dev。控制 ANALYZE 的度量指标详细程度。 |
TIMING | 布尔值 | METRICS 的语法糖:切换是否包含 timing 类别。 |
SUMMARY | 布尔值 | LEVEL 的语法糖:TRUE → summary,FALSE → dev。 |
COSTS | 布尔值 | 在普通 EXPLAIN 输出中包含统计信息(等价于 SET datafusion.explain.show_statistics)。不能与 ANALYZE 同时使用。 |
布尔参数可以省略不写(ANALYZE → true),也可以写成 TRUE/FALSE、ON/OFF 或 0/1。
与 ANALYZE 结合使用时,FORMAT 支持 indent(默认)和 pgjson;tree 和 graphviz 会被拒绝。pgjson 形式会输出带有实时度量指标的物理计划,这对于计划可视化工具非常有用——参见 ANALYZE 配合 pgjson 格式:
EXPLAIN (ANALYZE, FORMAT pgjson) SELECT ...;语句级选项优先于会话配置,因此你可以保留会话默认值不变,仅针对当前查询进行覆盖:
EXPLAIN (ANALYZE, LEVEL dev, METRICS 'rows,bytes') SELECT ...;DataFusion 未建模的 Postgres 选项(BUFFERS、WAL、SETTINGS、GENERIC_PLAN、MEMORY)会返回明确的错误,而不是被静默接受——请使用 METRICS 来过滤输出中显示的内容。
分区与执行
DataFusion 在查询规划过程中会确定要使用的最佳核心数。粗略地说,计划中的每个“分区”都会使用独立的核心分别执行。只有在 RepartitionExec、CoalescePartitions 和 SortPreservingMergeExec 等特定算子内部,数据才会在核心之间传递。
你可以在分区文档中了解更多相关信息。
聚合查询示例
让我们深入分析一个从 hits.parquet 文件聚合数据的查询示例。例如,ClickBench 中的这个查询可以按命中次数找出前 10 名用户:
SELECT "UserID", COUNT(*)
FROM 'hits.parquet'
GROUP BY "UserID"
ORDER BY COUNT(*) DESC
LIMIT 10;我们再次使用 EXPLAIN 查看查询计划:
> EXPLAIN FORMAT INDENT SELECT "UserID", COUNT(*) FROM 'hits.parquet' GROUP BY "UserID" ORDER BY COUNT(*) DESC LIMIT 10;
+---------------+---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
| plan_type | plan |
+---------------+---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
| logical_plan | Limit: skip=0, fetch=10 |
| | Sort: count(*) DESC NULLS FIRST, fetch=10 |
| | Aggregate: groupBy=[[hits.parquet.UserID]], aggr=[[count(Int64(1)) AS count(*)]] |
| | TableScan: hits.parquet projection=[UserID] |
| physical_plan | GlobalLimitExec: skip=0, fetch=10 |
| | SortPreservingMergeExec: [count(*)@1 DESC], fetch=10 |
| | SortExec: TopK(fetch=10), expr=[count(*)@1 DESC], preserve_partitioning=[true] |
| | AggregateExec: mode=FinalPartitioned, gby=[UserID@0 as UserID], aggr=[count(*)] |
| | CoalesceBatchesExec: target_batch_size=8192 |
| | RepartitionExec: partitioning=Hash([UserID@0], 10), input_partitions=10 |
| | AggregateExec: mode=Partial, gby=[UserID@0 as UserID], aggr=[count(*)] |
| | DataSourceExec: file_groups={10 groups: [[hits.parquet:0..1477997645], [hits.parquet:1477997645..2955995290], [hits.parquet:2955995290..4433992935], [hits.parquet:4433992935..5911990580], [hits.parquet:5911990580..7389988225], ...]}, projection=[UserID], file_type=parquet |
| | |
+---------------+---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+对于这个查询,我们再次从下往上读取计划:
逻辑计划算子
TableScanhits.parquet:从文件hits.parquet扫描数据。projection=[UserID]:只读取UserID列
AggregategroupBy=[[hits.parquet.UserID]]:按UserID列分组。aggr=[[count(Int64(1)) AS count(*)]]:对每个不同的分组应用COUNT聚合。
Sortcount(*) DESC NULLS FIRST:按计数降序排列数据。fetch=10:只返回前 10 行。
Limitskip=0:结果不跳过任何数据。fetch=10:将结果限制为 10 个值。
物理计划算子
DataSourceExecfile_groups={10 groups: [...]}:从hits.parquet文件并行读取 10 个分组。(上面的示例是在一台有 10 个核心的机器上运行的。)projection=[UserID]:下推UserID列的投影。Parquet 格式是列式的,DataFusion 的读取器只解码所需的列。
AggregateExecRepartitionExecpartitioning=Hash([UserID@0], 10):根据hash(UserID)的值将输入划分为 10 个(新的)输出分区。你可以在分区文档中了解更多信息。input_partitions=10:输入分区的数量。
CoalesceBatchesExectarget_batch_size=8192:将较小的批次合并成较大的批次。在这种情况下,每个批次大约包含 8192 行。
AggregateExecmode=FinalPartitioned:对每个分组执行最终聚合。更多信息请参阅多阶段分组的文档。gby=[UserID@0 as UserID]:按UserID分组。aggr=[count(*)]:对每个分组的所有行应用COUNT聚合。
SortExecTopK(fetch=10):使用一种特殊的 “TopK” 排序,一次只在内存中保留最大的 10 个值。你可以在 TopK 文档中了解更多信息。expr=[count(*)@1 DESC]:按降序对所有行进行排序。注意,这代表物理计划中的ORDER BY。preserve_partitioning=[true]:排序在每个分区上并行执行。在本例中,并行为 10 个分区分别找出前 10 个值。SortPreservingMergeExec[count(*)@1 DESC]:该操作符使用此表达式将 10 个不同的流合并为一个流。fetch=10:仅返回前 10 行
GlobalLimitExecskip=0:不跳过任何行fetch=10:仅返回前 10 行,即查询中的LIMIT 10。
本例中的数据
本节中的示例使用来自 ClickBench 的数据,这是一个数据分析基准测试。示例基于 14GB 的 hits.parquet 文件,可以从该网站下载,或使用以下命令下载:
cd benchmarks
./bench.sh data clickbench_1
***************************
DataFusion Benchmark Runner and Data Generator
COMMAND: data
BENCHMARK: clickbench_1
DATA_DIR: /Users/andrewlamb/Software/datafusion2/benchmarks/data
CARGO_COMMAND: cargo run --release
PREFER_HASH_JOIN: true
***************************
Checking hits.parquet...... found 14779976446 bytes ... Done然后你可以运行 datafusion-cli 来获取计划:
cd datafusion/benchmarks/data
datafusion-cli
DataFusion CLI v41.0.0
> select count(*) from 'hits.parquet';
+----------+
| count(*) |
+----------+
| 99997497 |
+----------+
1 row(s) fetched.
Elapsed 0.062 seconds.
>评论
登录后参与评论
KnowForge