
# SQL执行计划
SQL执行计划是一个节点树，显示DataArts Fabric SQL执行一条SQL语句时执行的详细步骤。
使用EXPLAIN命令可以查看优化器为每个查询生成的具体执行计划。EXPLAIN给每个执行节点都输出一行，显示基本的节点类型和优化器为执行这个节点预计的开销值。
#### 执行计划显示信息
除了设置不同的执行计划显示格式外，还可以通过不同的EXPLAIN用法，显示不同详细程度的执行计划信息。常见有如下几种，关于更多用法请参见[EXPLAIN](https://support.huaweicloud.com/devg-fabric/dataartsfabric_sql_04_0301.html)说明。
- EXPLAIN *statement*： 只生成执行计划，不实际执行。其中statement代表SQL语句。
- EXPLAIN ANALYZE *statement*：生成执行计划，进行执行，并显示执行的概要信息。显示中加入了实际的运行时间统计，包括在每个规划节点内部花掉的总时间(以毫秒计)和它实际返回的行数。
- EXPLAIN PERFORMANCE *statement*：生成执行计划，进行执行，并显示执行期间的全部信息。
为了测量运行时在执行计划中每个节点的开销，EXPLAIN ANALYZE或EXPLAIN PERFORMANCE会在当前查询执行上增加性能分析的开销。在一个查询上运行EXPLAIN ANALYZE或EXPLAIN PERFORMANCE有时会比普通查询明显地花费更多的时间。超支的数量依赖于查询的本质和使用的平台。
因此，当定位SQL运行慢问题时，如果SQL长时间运行未结束，建议通过EXPLAIN命令查看执行计划，进行初步定位。如果SQL可以运行出来，则推荐使用EXPLAIN ANALYZE或EXPLAIN PERFORMANCE查看执行计划及其实际的运行信息，以便更精准地定位问题原因。
**执行计划中的常见关键字说明：**
- 表访问方式 ForeignScan：全表顺序扫描。最基本的扫描算子，用于外表的顺序扫描，支持value分区剪枝。
  
- 表连接方式
  - Nested Loop 嵌套循环，适用于被连接的数据子集较小的查询。在嵌套循环中，外表驱动内表，外表返回的每一行都要在内表中检索找到它匹配的行，因此整个查询返回的结果集不能太大（不能大于10000），要把返回子集较小的表作为外表，而且在内表的连接字段上建议要有索引。
    
  
  - （Sonic）Hash Join 哈希连接，适用于数据量大的表的连接方式。优化器使用两个表中较小的表，利用连接键在内存中建立hash表，然后扫描较大的表并探测散列，找到与散列匹配的行。
    
   
- 运算符
  - sort 对结果集进行排序。
    
  
  - filter EXPLAIN输出显示。WHERE子句被用作一个"filter"条件，附加到顺序扫描的计划节点上。这意味着规划器会为该节点扫描的每一行检查该条件，并仅输出满足条件的行。虽然预计的输出行数因WHERE子句而减少，但顺序扫描仍需访问所有10000行。因此，总体开销并未降低，反而因检查WHERE条件而增加了额外的CPU时间（具体为10000 × cpu_operator_cost）。
    
  
  - LIMIT LIMIT限定了执行结果的输出记录数。如果增加了LIMIT，将不会检索到所有行。
    
   
 
#### 执行计划显示格式
DataArts Fabric SQL的执行计划显示格式层次清晰，计划包含了plan node id，性能分析简单直接。如[图1]。
图1pretty格式执行计划示例   
![](https://support.huaweicloud.com/devg-fabric/figure/zh-cn_image_0000002449670664.png "点击放大")
#### Stream计划
DataArts Fabric SQL中使用Stream计划进行查询的执行：
CN根据原语句生成计划并将计划下发给DN进行执行，各DN执行过程中使用Stream算子进行数据交互。
现有表tt01和tt02定义如下：
```
CREATE TABLE tt01(c1 int, c2 int)
store as orc;
CREATE TABLE tt02(c1 int, c2 int)
store as orc;
```
两表JOIN，且连接条件包含非分布列，其DN间存在数据交换。此时对于tt02表，会在各DN进行基表扫描，扫描后会通过Redistribute Stream算子，按照JOIN条件中的tt02.c1进行哈希计算后重新发送给各DN，然后在各DN上做JOIN，最后汇总到CN。
图2stream计划示例   
![](https://support.huaweicloud.com/devg-fabric/figure/zh-cn_image_0000002482790473.png "点击放大")
#### EXPLAIN PERFORMANCE详解
在SQL调优过程中经常需要执行EXPLAIN ANALYZE或EXPLAIN PERFORMANCE查看SQL语句实际执行信息，通过对比实际执行与优化器的估算之间的差别来为优化提供依据。EXPLAIN PERFORMANCE相对于EXPLAIN ANALYZE增加了每个DN上的执行信息。
以上一小节中的SQL查询语句为例：
```
SELECT * FROM tt01,tt02 WHERE tt01.c1=tt02.c1;
```
执行EXPLAIN PERFORMANCE输出的显示执行信息分为以下6个部分：
1. 执行计划
   图3explain performance显示信息示例   
   ![](https://support.huaweicloud.com/devg-fabric/figure/zh-cn_image_0000002449511028.png "点击放大")
   以表格的形式将计划显示出来，包含有11个字段，分别是：id、operation、A-time、A-rows、E-rows、E-distinct、Peak Memory、E-memory、A-width、E-width和E-costs。字段含义如下[表1]。
    表1执行字段说明 
   | 字段          | 描述                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                |
   |:---|:---|
   | id          | 执行算子节点编号。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                         |
   | operation   | 具体的执行节点算子名称。 Vector前缀的算子是指向量化执行引擎算子，一般出现含有列存表的Query中。 Streaming是一个特殊的算子，它实现了分布式架构的核心数据shuffle功能，Streaming共有三种形态，分别对应了分布式结构下不同的数据shuffle功能： - Streaming (type: GATHER)：作用是coordinator从DN收集数据。  - Streaming(type: REDISTRIBUTE)：作用是DN根据选定的列把数据重分布到所有的DN。  - Streaming(type: BROADCAST)：作用是把当前DN的数据广播给其他所有的DN。   |
   | A-time      | 各DN相应算子执行时间，一般DN上执行的算子的A-time是由\[\]括起来的两个值，分别表示此算子在所有DN上完成的最短时间和最长时间，包括下层算子执行时间。 注意：在整个计划中，除了叶子节点的执行时间是算子本身的执行时间，其余算子的执行时间均包含子节点的执行时间。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                   |
   | A-rows      | 表示相应算子输出的全局总行数。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                   |
   | E-rows      | 每个算子估算的输出行数。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                      |
   | E-distinct  | 表示hashjoin算子的distinct估计值。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                         |
   | Peak Memory | 此算子在每个DN上执行时使用的内存峰值，\[\]中左侧为最小值，右侧为最大值。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                           |
   | E-memory    | DN上每个算子估算的内存使用量，只有DN上执行的算子会显示。某些场景会在估算的内存使用量后使用括号显示该算子在内存资源充足下可以自动扩展的内存上限。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                        |
   | A-width     | 表示当前算子每行元组的实际宽度，仅对于重内存使用算子会显示，包括：(Vec)HashJoin、(Vec)HashAgg、(Vec) HashSetOp、(Vec)Sort、(Vec)Materialize算子等，其中(Vec)HashJoin计算的宽度是其右子树算子的宽度，会显示在其右子树上。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                               |
   | E-width     | 每个算子输出元组的估算宽度。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                    |
   | E-costs     | 每个算子估算的执行代价。 - E-costs是优化器根据成本参数定义的单位来衡量的，习惯上以磁盘页面抓取为1个单位， 其它开销参数将参照它来设置。  - 每个节点的开销（E-costs值）包括它的所有子节点的开销。  - 开销只反映了优化器关心的东西，并没有把结果行传递给客户端的时间考虑进去。虽然这个时间可能在实际的总时间里占据相当重要的分量，但是被优化器忽略了，因为它无法通过修改规划来改变。                                                                                                                                                                                                                                                                                                                       |
      
   
2. Predicate Information (identified by plan id) ![](https://support.huaweicloud.com/devg-fabric/figure/zh-cn_image_0000002482710453.png "点击放大")
   谓词过滤这部分主要显示的是对应执行算子节点的过滤条件，即在整个计划执行过程中不会变的信息，主要是一些join条件和一些filter信息。对于分区表，还会显示分区剪枝信息。
   
3. Memory Information (identified by plan id) ![](https://support.huaweicloud.com/devg-fabric/figure/zh-cn_image_0000002482790521.png "点击放大")
   内存使用信息这部分显示的是整个计划中会将内存的使用情况打印出来的算子的内存使用信息，主要是Hash、Sort算子，包括算子峰值内存（peak memory）、优化器预估的内存（estimate memory）、控制内存（control memory）、估算内存使用（operator memory）、执行时实际宽度（width）、内存使用自动扩展次数（auto spread num）、是否提前下盘（early spilled）以及下盘信息，包括重复下盘次数（spill Time(s)）、内外表下盘分区数（inner/outer partition spill num）、下盘文件数（temp file num）、下盘数据量及最小和最大分区的下盘数据量（written disk IO \[min, max\] ）。其中sort算子不会显示具体的下盘文件数，仅在显示排序方法时显示Disk。下方是一个发生了下盘的算子的内存信息示例。
   ![](https://support.huaweicloud.com/devg-fabric/figure/zh-cn_image_0000002449511008.png "点击放大")
   
4. Targetlist Information (identified by plan id) ![](https://support.huaweicloud.com/devg-fabric/figure/zh-cn_image_0000002482710461.png "点击放大")
   这一部分显示的是每一个算子对应的输出目标列信息。
   
5. DataNode Information (identified by plan id) ![](https://support.huaweicloud.com/devg-fabric/figure/zh-cn_image_0000002449510964.png "点击放大")
   这部分将各个算子的执行时间（如果包含过滤及投影也会显示对应的执行时间）、CPU、buffer的使用情况全部打印出来。
   - 算子执行信息 ![](https://support.huaweicloud.com/devg-fabric/figure/zh-cn_image_0000002482790529.png "点击放大")
     每个算子的执行信息都包含三个部分：
     - executor0_es_group/executor1_es_group表示具体执行的节点信息，括号中的信息是实际的执行信息。
     
     - actual time表示实际的执行时间，第一个数字表示执行时进入当前算子到输出第一条数据所花费的时间，第二个数字表示输出所有数据的总执行时间。
     
     - rows表示当前算子输出数据行数。
     
     
     
     - loops表示当前算子的执行次数。需要注意，对于分区表来说，每一个分区表的扫描就是一次完整的扫描操作，当切换到下一个分区的时候，又是一次新的扫描操作。
       
   
   - CPU信息 ![](https://support.huaweicloud.com/devg-fabric/figure/zh-cn_image_0000002449670612.png "点击放大")
     每个算子执行的过程都有CPU信息，其中cyc代表的是CPU的周期数，ex cyc表示的是当前算子的周期数，不包含其子节点；inc cyc是包含子节点的周期数；ex row是当前算子输出的数据行数；ex c/r则是ex cyc/ex row得到的每条数据所用的平均周期数。
     
   
   - Buffer信息 ![](https://support.huaweicloud.com/devg-fabric/figure/zh-cn_image_0000002482710481.png "点击放大")
     buffers显示缓冲区信息，包括共享块和临时块的读和写。
     共享块包含表和索引，临时块在排序和物化中使用的磁盘块。上层节点显示出来的块数据包含了其所有子节点使用的块数。
     对于发生了下盘的算子，Buffer信息会展示下盘的数据量，"temp read"代表读取下盘的临时数据的次数，"written"代表写下盘的临时数据的次数，"written_size"代表下盘的数据量。
     当关闭spill-to-OBS特性GUC开关enable_spill_to_remote_storage时，数据下盘到DN实例目录中，数据显示格式如下：
     ![](https://support.huaweicloud.com/devg-fabric/figure/zh-cn_image_0000002449510992.png "点击放大")
     当打开spill-to-OBS特性GUC开关enable_spill_to_remote_storage时，还会额外显示和该特性相关的下盘统计信息，示例如下图所示。其中，written_disk_size代表下盘到磁盘缓存的数据量；written_obs_size代表如磁盘缓存空间不足情况时直接下盘到obs的数据量；spill_obs_size代表的是磁盘缓存空间不足时，从磁盘缓存写回到obs的数据量。
     ![](https://support.huaweicloud.com/devg-fabric/figure/zh-cn_image_0000002482790485.png "点击放大")
     
   
   - 磁盘缓存信息 ![](https://support.huaweicloud.com/devg-fabric/figure/zh-cn_image_0000002449670648.png "点击放大")
     - Disk Cache：表示磁盘缓存的命中信息和数据读取信息。
       - miss代表磁盘缓存未命中的次数。
       
       - hit表示磁盘缓存命中的次数。
       
       - disk_cache_error_code，error表示产生errorCode的次数。
       
       - scanBytes表示scan查询的数据量。
       
       - remoteReadBytes表示在OBS上读取的数据量。
       
       - loadTime表示从磁盘缓存加载数据的时间。
       
       - openTime表示打开磁盘缓存文件的时间。
       
       - preadTime表示从磁盘读取数据的时间。 注：为了提升OBS的效率会对相邻的请求块合并，或者因为请求写磁盘缓存的最小粒度是block(默认1M)，可能会使得scanBytes会小于remoteReadBytes。
         
        
     
     
     
     - ReadAhead：表示数据预取的相关信息。
       - parseMetaTime表示预读时解析文件元数据的时间。
       
       - submitBytes表示预读的数据量。
       
       - submitTime表示数据预读请求的提交时间。
       
       - waitTime表示读取数据时命中执行中的数据预读请求并且等待预读请求完成的时间。
       
       - waitCount表示读取数据时命中执行中的数据预读请求的次数。
       
       - cancelCount表示读取数据时命中并撤销尚未开始执行的数据预读请求的次数。
       
       - hitCount表示读取数据时命中的已经完成的数据预读请求的次数。
       
       - fabricCacheHitCount L1Cache表示数据缓存在L1层命中的次数。
       
       - L2Cache表示数据缓存在L2层命中的次数。
        
     
     
     
     - OBS I/O ：表示OBS IO请求的详细信息。
       1. count表示OBS IO请求的总数量。
       
       2. averageRTT表示OBS IO请求的平均RTT(Round Trip Time，IO请求往返时间)，单位为μs。
       
       3. averageLatency表示OBS IO请求的平均延迟，单位为μs。
       
       4. latencyGt1s表示OBS IO请求延迟超过1s的数量。
       
       5. latencyGt10s表示OBS IO请求延迟超过10s的数量。
       
       6. retryCount表示OBS IO请求重试的总次数。
       
       7. rateLimitCount表示OBS IO请求被流控的总次数。
        
      
   
   - 元戎接口调用信息 ![](https://support.huaweicloud.com/devg-fabric/figure/zh-cn_image_0000002482710525.png "点击放大")
     列出了所有executor actor上元戎接口调用统计信息，并在最后给出了总体信息。
     每一行表示一个接口的统计信息，其字段的含义请参考函数yr_api_status()的字段。
     
   
   - Query Summary ![](https://support.huaweicloud.com/devg-fabric/figure/zh-cn_image_0000002449510956.png "点击放大")
     这一部分主要打印语句级执行信息，包括了各阶段执行时间、各个DN上初始化和结束阶段的最大最小执行时间、CN上的初始化、执行、结束阶段的时间，以及当前语句执行时系统可用内存、语句估算内存、并行度等信息。
     - DN执行信息
       - DataNode executor start time：DN执行器开始时间，\[min_node_name, max_node_name\] : \[min_time, max_time\]。
       
       - DataNode executor run time：DN执行器运行时间，\[min_node_name, max_node_name\] : \[min_time, max_time\]。
       
       - DataNode executor end time：DN执行器结束时间，\[min_node_name, max_node_name\] : \[min_time, max_time\]。
        
     
     - Remote query poll time：接收结果时用于poll等待的时间。
     
     - 内存估算信息
       - System available mem：系统可用内存。
       
       - Query Max mem：查询最大内存。
       
       - Query estimated mem：语句估算内存。
        
     
     - SMP自适应并行信息（仅开启SMP自适应时显示）
       - Initial DOP：计划生成的初始规划并行度。
       
       - Avail max core：可供该语句执行使用的CPU核数。
       
       - Final Max DOP：计划中算子的最大并行度。
        
     
     - lakeformation访问信息：包括信息访问时间及次数。
     
     - CN执行时间
       - Coordinator executor start time：CN执行器开始时间。
       
       - Coordinator executor run time：CN执行器运行时间。
       
       - Coordinator executor end time：CN执行器结束时间。
        
     
     - Parser runtime：解析器运行时间。
     
     - YR utilize function analyse：元戎调用函数分析，包括：
       - Actor invoke runtime：拉起actor的时间。
       
       - Actor Distribution Info：拉起actor的IP及序号信息。
       
       - Actor Stream topo create time：actor的stream启动时间。
       
       - Actor Consumer close time：actor的consumer关闭时间。
       
       - Actor Stream Producer/consumer count：actor的stream个数。
        
     
     - Planner runtime：优化器执行时间。
     
     - Query Id：查询ID。
     
     - Total runtime：总执行时间。
      
    
![](https://support.huaweicloud.com/devg-fabric/public_sys-resources/notice_3.0-zh-cn.png)
- A-rows和E-rows的差异体现了优化器估算和实际执行的偏差度。一般情况下两者偏差越大，则可以认为优化器生成的计划的越不可信，人工干预调优的必要性越大。
- A-time中的两个值偏差越大，表明此算子的计算偏斜(在不同DN上执行时间差异)越大，人工干预调优的必要性越大。一般来说，两个相邻的算子，上层算子的执行时间包含下层算子的执行时间，但如果上层算子为stream算子，由于各线程不存在驱动关系，上层算子执行时间可能小于下层算子的执行时间，即不存在包含关系。
- Max Query Peak Memory经常用来估算SQL语句耗费内存，也被用来作为SQL语句调优时运行态内存参数设置的重要依据。一般会以EXPLAIN ANALYZE或EXPLAIN PERFORMANCE的输出作为进一步调优的输入。
 
