Operators

Updated at:

This topic describes common operators in execution plans and explains what each operator does.

Operator reference

Category

Operators

Push down to data nodes

LogicalView, LogicalModifyView, PhyTableOperation, IndexScan

Join

BKAJoin, NLJoin, HashJoin, SortMergeJoin, HashSemiJoin, SortMergeSemiJoin, MaterializedSemiJoin

Sort

MemSort, TopN, MergeSort

Aggregate (GROUP BY)

HashAgg, SortAgg

Redistribute or collect data

Exchange, Gather

Filter

Filter

Project

Project

Merge result sets

UnionAll, UnionDistinct

Limit output rows (Limit/Offset...Fetch)

Limit

Window function

OverWindow

Operators that push down to data nodes

LogicalView

LogicalView reads data from the storage-layer MySQL data source. It is similar to TableScan or IndexScan in other databases but supports a broader set of pushdown operations. LogicalView contains the pushed-down SQL statement and data source information, and behaves more like a view. The pushed-down SQL may include Project, Filter, aggregation, sort, Join, and subquery operations. The following example shows the output of LogicalView and explains each field:

explain select * From sbtest1 where id > 1000;

Output:

Gather(concurrent=true)
   LogicalView(tables="[0000-0031].sbtest1_[000-127]", shardCount=128, sql="SELECT * FROM `sbtest1` WHERE (`id` > ?)")

LogicalView consists of three parts:

  • tables: the table name in the storage-layer MySQL, split by a period.. The part before the period. is the database shard number, and the part after is the table name with its shard number. For example, [000-127] represents all tables with shard numbers from 000 to 127.

  • shardCount: the total number of table shards to access. In this example, 128 table shards from 000 to 127 are accessed.

  • sql: the SQL template sent to the storage-layer MySQL. PolarDB-X At runtime, the table name is replaced with the physical table name, and the parameterized constant placeholder? is replaced with the actual parameter value. For more information, see Execution plan management.

LogicalModifyView

LogicalModifyView modifies data in the underlying data source. It also records an SQL statement, which can be INSERT, UPDATE, or DELETE. The following examples show the output of LogicalModifyView and explain each field:

  • Example 1

    explain update sbtest1 set c='Hello, DRDS' where id > 1000;

    Output:

    LogicalModifyView(tables="[0000-0031].sbtest1_[000-127]", shardCount=128, sql="UPDATE `sbtest1` SET `c` = ? WHERE (`id` > ?)"
  • Example 2

    explain delete from sbtest1 where id > 1000;

    Output:

    LogicalModifyView(tables="[0000-0031].sbtest1_[000-127]", shardCount=128, sql="DELETE FROM `sbtest1` WHERE (`id` > ?)")

The fields in the LogicalModifyView query plan are similar to those in LogicalView, including the physical table shards, shard count, and SQL template. Similarly, because the execution plan cache is enabled, the SQL is parameterized and constants in the SQL template are replaced with?.

PhyTableOperation

PhyTableOperation performs an operation directly on a single physical table shard.

Note

Typically, this operator is used only for INSERT statements. However, when a query is routed to a single shard, this operator may also appear in SELECT statements.

explain insert into sbtest1 values(1, 1, '1', '1'),(2, 2, '2', '2');

Output:

PhyTableOperation(tables="SYSBENCH_CORONADB_1526954857179TGMMSYSBENCH_CORONADB_VGOC_0000_RDS.[sbtest1_001]", sql="INSERT INTO ? (`id`, `k`, `c`, `pad`) VALUES(?, ?, ?, ?)", params="`sbtest1_001`,1,1,1,1")
PhyTableOperation(tables="SYSBENCH_CORONADB_1526954857179TGMMSYSBENCH_CORONADB_VGOC_0000_RDS.[sbtest1_002]", sql="INSERT INTO ? (`id`, `k`, `c`, `pad`) VALUES(?, ?, ?, ?)", params="`sbtest1_002`,2,2,2,2")

In this example, the INSERT statement inserts two rows, and each row corresponds to one PhyTableOperation operator. PhyTableOperation consists of three parts:

  • tables: the physical table name. Each PhyTableOperation targets exactly one physical table.

  • sql: the SQL template, in which the table name and constants are parameterized and replaced with?. The corresponding parameter values are provided in the params field that follows.

  • params: the parameter values for the SQL template, including the table name and constants.

IndexScan

IndexScan reads data from the storage-layer MySQL data source by scanning an index table, similar to LogicalView. The following example shows the output of IndexScan and explains each field:

explain select * from sequence_one_base where integer_test=1;

Output:

+-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
| IndexScan(tables="DRDS_POLARX1_QATEST_APP_000000_GROUP.gsi_sequence_one_index_3a0A_01", sql="SELECT `pk`, `integer_test`, `varchar_test`, `char_test`, `blob_test`, `tinyint_test`, `tinyint_1bit_test`, `smallint_test`, `mediumint_test`, `bit_test`, `bigint_test`, `float_test`, `double_test`, `decimal_test`, `date_test`, `time_test`, `datetime_test`, `timestamp_test`, `year_test`, `mediumtext_test` FROM `gsi_dml_sequence_one_index_index1` AS `gsi_dml_sequence_one_index_index1` WHERE (`integer_test` = ?)") |

The preceding SQL statement would normally scan the sequence_one_base table. Because integer_test is not a partition key, scanning all shards of sequence_one_base would be required. However, because the sequence_one_base table has a global secondary index (GSI) named gsi_sequence_one_index on the integer_test column, the predicate integer_test=1 is used to prune the index table. As a result, an IndexScan operator is generated, which scans only one shard of the gsi_sequence_one_index table.

Operators that run on compute nodes or data nodes

UnionAll and UnionDistinct

UnionAll corresponds to UNION ALL, and UnionDistinct corresponds to UNION DISTINCT. This operator typically has two or more inputs and merges data from multiple inputs into one. Example:

explain select * From sbtest1 where id > 1000 union distinct select * From sbtest1 where id < 200;

Output:

UnionDistinct(concurrent=true)
  Gather(concurrent=true)
    LogicalView(tables="[0000-0031].sbtest1_[000-127]", shardCount=128, sql="SELECT * FROM `sbtest1` WHERE (`id` > ?)")
  Gather(concurrent=true)
    LogicalView(tables="[0000-0031].sbtest1_[000-127]", shardCount=128, sql="SELECT * FROM `sbtest1` WHERE (`id` < ?)")

Gather

Gather merges multiple data streams into a single data stream. In the preceding example, Gather merges data collected from each table shard into one result set. Gather typically appears above LogicalView and collects data from all scanned shards.

Exchange

Exchange is a logical operator that does not perform any computation on the data. It redistributes the input data and passes it to downstream operators. Common redistribution strategies include:

  • SINGLETON: merges multiple upstream data streams into one output stream. This strategy is equivalent to Gather.

  • HASH_DISTRIBUTED: repartitions the upstream input data by specific columns. This strategy is commonly used in execution plans that contain Join and Agg operators.

  • BROADCAST_DISTRIBUTED: distributes a single copy of the upstream data to multiple downstream nodes. This strategy is primarily used in MPP execution plans.

MergeSort

MergeSort is a merge sort operator that merges multiple sorted data streams into a single sorted data stream. Example:

explain select * from sbtest1 where id > 1000 order by id limit 5,10; 

Output:

MergeSort(sort="id ASC", offset=?1, fetch=?2)   
   LogicalView(tables="[0000-0031].sbtest1_[000-127]", shardCount=128, sql="SELECT * FROM `sbtest1` WHERE (`id` > ?) ORDER BY `id` LIMIT (? + ?)")

MergeSort consists of three parts:

  • sort: the sort field and sort order. id ASC means ascending order by the id field. DESC means descending order.

  • offset: the number of rows to skip when retrieving the result set. The value is parameterized in the example; the actual value is 5.

  • fetch: the maximum number of rows to return. Similar to offset, the value is parameterized; the actual value is 10.

Project

Project performs a projection operation, which selects specific columns from the input data for output, or transforms certain columns through functions or expressions before output. It can also include constants.

explain select 'hello, DRDS', 1 / 2, CURTIME(); 

Output:

Project(hello, DRDS="_UTF-16'hello, DRDS'", 1 / 2="1 / 2", CURTIME()="CURTIME()")

The Project plan includes each output column name alongside its corresponding column, value, function, or expression.

Filter

Filter applies filter conditions to the input data. Rows that satisfy the conditions are passed through; otherwise, they are discarded. The following is a complex example that includes most of the operators described above.

explain select k, avg(id) avg_id from sbtest1 where id > 1000 group by k having avg_id > 1300;

Output:

Filter(condition="avg_id > ?1")
  Project(k="k", avg_id="sum_pushed_sum / sum_pushed_count")
    SortAgg(group="k", sum_pushed_sum="SUM(pushed_sum)", sum_pushed_count="SUM(pushed_count)")
      MergeSort(sort="k ASC")
        LogicalView(tables="[0000-0031].sbtest1_[000-127]", shardCount=128, sql="SELECT `k`, SUM(`id`) AS `pushed_sum`, COUNT(`id`) AS `pushed_count` FROM `sbtest1` WHERE (`id` > ?) GROUP BY `k` ORDER BY `k`")

The condition WHERE id>1000 does not have a corresponding Filter operator because it is pushed down into LogicalView. You can see WHERE (id > ?) in the SQL template of LogicalView.