Чтение планов запросов SQL-движка на основе Apache Calcite#

План запроса для SQL-движка на основе Apache Calcite представляет собой дерево физических реляционных операторов. Источником данных являются листья этого дерева. Данные проходят снизу вверх, от листьев к корню, и на выходе из корневого элемента получается набор данных для пользователя. Каждый узел дерева, помимо информации о типе реляционного оператора, также включает в себя атрибуты этого оператора, характеристики оператора и оценочную информацию о стоимости выполнения ветки дерева начиная с данного узла.

Атрибуты оператора#

Набор атрибутов специфичен для каждого типа оператора. Атрибуты и их значения в плане отображаются сразу после типа оператора, например:

IgniteTableScan(table=[[PUBLIC, EMPLOYER]], requiredColumns=[{2, 3}])

В данном случае у оператора есть два явно заданных атрибута: table и requiredColumns.

Характеристики оператора#

Набор характеристик является общим для всех типов операторов. Характеристики не выводятся в плане, но каждый узел содержит их в неявном виде. Характеристики обычно наследуются оператором от нижележащих операторов — но есть и операторы, которые модифицируют характеристики.

Основные характеристики:

  • COLLATION (упорядоченность данных) определяет, по каким колонкам и в каком направлении (по возрастанию или по убыванию) отсортированы полученные от оператора данные. Неявным образом эта характеристика возникает при сканировании индекса (данные отсортированы по полям индекса), явным образом может быть переопределена с помощью оператора сортировки. Некоторые операторы нарушают упорядоченность данных и, соответственно, убирают эту характеристику. Например, после хеш-агрегации данные не будут отсортированы.

  • REWINDABILITY (возможность обратной перемотки данных) — практически все типы операторов (за исключением типа оператора «обмен» — IgniteExchange) могут повторно итерироваться по данным (циклично их перебирать), если входные данные также способны к обратной перемотке. Оператор «обмен» может итерироваться по данным только один раз. Если для данных, которые не способны перематываться, требуется возможность повторного итерирования (например, для типа оператора «коррелированное объединение» — Correlated Nested Loop Join, подробнее о нем написано ниже в разделе «IgniteCorrelatedNestedLoopJoin»), для изменения этой характеристики используется тип оператора Spool. Он кеширует все входные данные и обладает возможностью повторно по ним итерироваться.

  • DISTRIBUTION (распределение данных по кластеру) может принимать значения:

    • single — данные присутствуют на узле-инициаторе запроса;

    • broadcast — данные присутствуют на всех узлах (REPLICATED-кеши, VALUES);

    • hash — данные распределены между узлами кластера по хеш-функции от перечисленных полей;

    • affinity — частный случай хеш-распределения, данные распределены по аффинити-функции (PARTITIONED-кеши);

    • random — данные распределены случайным образом.

    Для изменения распределения данных используется оператор «обмен» — IgniteExchange, подробнее о нем написано ниже в разделе «IgniteExchange». По операторам IgniteExchange идет разделение плана на фрагменты. Количество фрагментов плана = количество IgniteExchange + 1. Каждый фрагмент может отправляться на свой набор узлов кластера и пересылать данные вышестоящему фрагменту также на набор узлов кластера, на который был отправлен вышестоящий фрагмент. У корневого фрагмента (и корневого узла дерева) всегда будет single-распределение, то есть он всегда выполняется на узле-инициаторе запроса (результат запроса отдается пользователю с помощью узла-инициатора запроса). Если в плане нет операторов IgniteExchange, у плана будет только один фрагмент, который будет выполняться только на узле-инициаторе запроса.

Стоимость оператора#

Стоимость ветки (текущего оператора и всех нижележащих) оценивается на основе статистики и предположений о селективности данных. Она может быть далека от реальности и нужна только для сравнения между собой различных планов и выбора оптимального. В стоимость входят:

  • количество строк на выходе из оператора;

  • общее количество строк, которые обработаны в ветке плана;

  • затраты на CPU, память и сеть в условных единицах.

Пример, как может выглядеть информация о стоимости в плане:

rowcount = 1.0, cumulative cost = IgniteCost [rowCount=19.0, cpu=28.0, memory=9.0, io=0.0, network=36.0]

Типы операторов#

Источники данных (листовые узлы)#

IgniteValues#

Значения, которые статически заданы в запросе.

Атрибуты:

  • tuples — массив выводимых строк.

SELECT * FROM (VALUES (0, 'ROW0'), (1, 'ROW1'))
IgniteValues(tuples=[[{ 0, _UTF-8'ROW0' }, { 1, _UTF-8'ROW1' }]]): rowcount = 2.0, cumulative cost = IgniteCost [rowCount=2.0, cpu=1.0, memory=0.0, io=0.0, network=0.0], id = xx

IgniteTableScan#

Сканирование таблицы.

Атрибуты:

  • table — имя таблицы;

  • filters — предикат;

  • project — проекция;

  • requiredColumns — список необходимых колонок (если отсутствует, будут отсканированы все колонки).

SELECT TO_CHAR(ORDER_DATE, 'DD.MM.YYYY') FROM ORDERS WHERE TOTAL_SUM > 1000
IgniteExchange(distribution=[single]): rowcount = 250000.0, cumulative cost = IgniteCost [rowCount=750000.0, cpu=2250000.0, memory=0.0, io=0.0, network=1000000.0], id = xx
  IgniteTableScan(table=[[PUBLIC, ORDERS]], filters=[>($t1, 1000)], projects=[[TO_CHAR(CAST($t0):TIMESTAMP(3) NOT NULL, _UTF-8'DD.MM.YYYY')]], requiredColumns=[{1, 2}]): rowcount = 250000.0, cumulative cost = IgniteCost [rowCount=500000.0, cpu=2000000.0, memory=0.0, io=0.0, network=0.0], id = xx

IgniteIndexScan#

Сканирование индекса.

Атрибуты:

  • table — имя таблицы;

  • index — имя индекса;

  • searchBounds — границы сканирования индекса;

  • filters — предикат (может быть более жестким, чем границы сканирования индекса);

  • project — проекция;

  • inlineScan — возможность использования индекса только по встроенным (inline) данным в индексе без обращения к страницам данных;

  • requiredColumns — список необходимых колонок (если отсутствует, будут отсканированы все колонки);

  • collation — порядок данных в индексе.

SELECT TO_CHAR(ORDER_DATE, 'DD.MM.YYYY') FROM ORDERS WHERE ORDER_ID < 10 AND TOTAL_SUM > 1000
IgniteExchange(distribution=[single]): rowcount = 125000.0, cumulative cost = IgniteCost [rowCount=375001.0, cpu=1125040.367090132, memory=1.0, io=1.0, network=500001.0], id = xx
  IgniteIndexScan(table=[[PUBLIC, ORDERS]], index=[ORDER_ID_IDX], filters=[AND(<($t0, 10), >($t2, 1000))], projects=[[TO_CHAR(CAST($t1):TIMESTAMP(3) NOT NULL, _UTF-8'DD.MM.YYYY')]], requiredColumns=[{0, 1, 2}], searchBounds=[[RangeBounds [lowerBound=$NULL_BOUND(), upperBound=10, lowerInclude=false, upperInclude=false], null, null]], inlineScan=[false], collation=[[0 ASC-nulls-first]]): rowcount = 125000.0, cumulative cost = IgniteCost [rowCount=250001.0, cpu=1000040.3670901322, memory=1.0, io=1.0, network=1.0], id = xx

IgniteTableFunctionScan#

Сканирование табличной функции.

Атрибуты:

  • invocation — функция.

SELECT * FROM table(system_range(1, 10))
IgniteTableFunctionScan(invocation=[SYSTEM_RANGE(1, 10)], rowType=[RecordType(BIGINT X)]): rowcount = 100.0, cumulative cost = IgniteCost [rowCount=100.0, cpu=100.0, memory=0.0, io=0.0, network=0.0], id = xx

IgniteIndexCount#

Подсчет количества строк по индексу без обращения к страницам данных. Подсчет возможен, только если нет других фильтров и группировки или если по индексированной колонке используются COUNT(*) или COUNT(COL).

Атрибуты:

  • table — имя таблицы;

  • index — имя индекса;

  • notNull — требуются только непустые значения (COUNT(COL));

  • fieldIdx — порядковый номер нужной колонки в индексе.

SELECT COUNT(ORDER_ID) FROM ORDERS
IgniteProject($f0=[CAST($0):BIGINT NOT NULL]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=4.0, cpu=503.0, memory=5.0, io=0.0, network=4.0], id = xx
  IgniteColocatedHashAggregate(group=[{}], agg#0=[$SUM0($0)]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=3.0, cpu=502.0, memory=5.0, io=0.0, network=4.0], id = xx
    IgniteExchange(distribution=[single]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=2.0, cpu=501.0, memory=0.0, io=0.0, network=4.0], id = xx
      IgniteIndexCount(index=[ORDER_ID_IDX], table=[[PUBLIC, ORDERS]], notNull=[true], fieldIdx=[0]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=1.0, cpu=500.0, memory=0.0, io=0.0, network=0.0], id = xx

IgniteIndexBound#

Определение минимального/максимального значения по индексу без итерирования по индексу или данным. Определение возможно, только если нет других фильтров и используется MIN(COL)/MAX(COL) по индексированной колонке (если эта колонка первая в индексе).

Атрибуты:

  • table — имя таблицы;

  • index — имя индекса;

  • first — требуются первое (true, MIN) или последнее (false, MAX) значение;

  • fieldIdx — порядковый номер нужной колонки в индексе;

  • requiredColumns — список необходимых колонок (если отсутствует, будут использоваться все колонки).

SELECT MIN(ORDER_ID) FROM ORDERS
IgniteReduceHashAggregate(rowType=[RecordType(INTEGER EXPR$0)], group=[{}], EXPR$0=[MIN($0)]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=4.0, cpu=25003.00037252903, memory=10.0, io=0.0, network=12.0], id = xx
  IgniteExchange(distribution=[single]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=3.0, cpu=25002.00037252903, memory=5.0, io=0.0, network=12.0], id = xx
    IgniteMapHashAggregate(group=[{}], EXPR$0=[MIN($0)]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=2.0, cpu=25001.00037252903, memory=5.0, io=0.0, network=0.0], id = xx
      IgniteIndexBound(table=[[PUBLIC, ORDERS]], index=[ORDER_ID_IDX], first=[true], requiredCols=[{0}]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=1.0, cpu=25000.00037252903, memory=0.0, io=0.0, network=0.0], id = xx

Обмены между узлами кластера#

IgniteExchange#

Приведение одного распределения данных в кластере к другому с помощью рассылки по определенным правилам данных между узлами кластера. Обычно данные приводятся к распределению в виде:

  • single — для случаев, когда требуется дальнейшая обработка на единственном узле, например, для reduce-фазы агрегатов/операций над множествами или для обработки финального результата на узле-инициаторе с последующей передачей клиенту;

  • hash/affinity — для коллокации данных, например, в случае соединения изначально не коллоцированных таблиц.

Для данных, которые уже присутствуют на всех узлах кластера (broadcast), рассылка не нужна. Вместо этого используется оператор IgniteTrimExchange (подробнее о нем написано ниже в разделе «IgniteTrimExchange»), который оставляет в выборке только данные, соответствующие новому распределению.

Атрибуты:

  • distribution — распределение, к которому требуется привести.

Пример запроса, в котором данные в таблице приведены к affinity-распределению:

SELECT * FROM ORDERS
IgniteExchange(distribution=[single]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=1000000.0, cpu=1000000.0, memory=0.0, io=0.0, network=6000000.0], id = xx
  IgniteTableScan(table=[[PUBLIC, ORDERS]], requiredColumns=[{0, 1, 2}]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500000.0, cpu=500000.0, memory=0.0, io=0.0, network=0.0], id = xx

Пояснения к плану: у данных в таблице ORDERS распределение affinity (ORDER_ID, ORDERS_CACHE, RendezvousAffinityFunction). Они находятся на разных узлах кластера, и чтобы отдать их клиенту, нужно переместить их на узел-инициатор запроса (distribution=[single]).

Пример запроса, в котором данные в таблице расположены на разных узлах кластера:

SELECT AVG(TOTAL_SUM) FROM ORDERS
IgniteReduceHashAggregate(rowType=[RecordType(DECIMAL(32767, 0) EXPR$0)], group=[{}], EXPR$0=[AVG($0)]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=1000002.0, cpu=1000002.0, memory=10.0, io=0.0, network=12.0], id = xx
  IgniteExchange(distribution=[single]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=1000001.0, cpu=1000001.0, memory=5.0, io=0.0, network=12.0], id = xx
    IgniteMapHashAggregate(group=[{}], EXPR$0=[AVG($0)]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=1000000.0, cpu=1000000.0, memory=5.0, io=0.0, network=0.0], id = xx
      IgniteTableScan(table=[[PUBLIC, ORDERS]], requiredColumns=[{2}]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500000.0, cpu=500000.0, memory=0.0, io=0.0, network=0.0], id = xx

Пояснения к плану: данные в таблице ORDERS расположены на разных узлах кластера. В этом случае у выборки random-распределение, так как в ней отсутствует аффинити-ключ, по которому распределены данные. В любом случае (даже если бы результат агрегации не передавался напрямую пользователю, а участвовал в дальнейшей обработке) чтобы выполнить reduce-фазу агрегата при random-, hash и affinity-распределениях, нужно переместить данные на один узел.

Пример запроса, в котором таблицы не коллоцированы:

SELECT o.order_date, oi.item_id FROM orders o JOIN order_items oi ON (o.order_id = oi.order_id) WHERE oi.amount = 1
IgniteExchange(distribution=[single]): rowcount = 11250.0, cumulative cost = IgniteCost [rowCount=1672502.0, cpu=4897502.0, memory=2.0, io=2.0, network=990002.0], id = xxxxx
  IgniteProject(ORDER_DATE=[$1], ITEM_ID=[$3]): rowcount = 11250.0, cumulative cost = IgniteCost [rowCount=1661252.0, cpu=4886252.0, memory=2.0, io=2.0, network=900002.0], id = xxxxx
    IgniteMergeJoin(condition=[=($0, $2)], joinType=[inner], variablesSet=[[]], leftCollation=[[0 ASC-nulls-first]], rightCollation=[[0 ASC-nulls-first]]): rowcount = 11250.0, cumulative cost = IgniteCost [rowCount=1650002.0, cpu=4875002.0, memory=2.0, io=2.0, network=900002.0], id = xxxxx
      IgniteIndexScan(table=[[PUBLIC, ORDERS]], index=[ORDER_ID_IDX], requiredColumns=[{0, 1}], inlineScan=[false], collation=[[0 ASC-nulls-first]]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500001.0, cpu=500001.0, memory=1.0, io=1.0, network=1.0], id = xxx
      IgniteExchange(distribution=[affinity[identity=org.apache.ignite.cache.affinity.rendezvous.RendezvousAffinityFunction, cacheId=xxxxxxxxxx][0]]): rowcount = 75000.0, cumulative cost = IgniteCost [rowCount=575001.0, cpu=2075001.0, memory=1.0, io=1.0, network=900001.0], id = xxxxx
        IgniteIndexScan(table=[[PUBLIC, ORDER_ITEMS]], index=[ORDER_ITEMS_ORDER_ID_IDX], filters=[=($t2, 1)], requiredColumns=[{1, 2, 3}], inlineScan=[false], collation=[[1 ASC-nulls-first]]): rowcount = 75000.0, cumulative cost = IgniteCost [rowCount=500001.0, cpu=2000001.0, memory=1.0, io=1.0, network=1.0], id = xxx

Пояснения к плану: в данном примере таблицы ORDERS и ORDER_ITEMS не коллоцированы. Соединение таблиц по колонкам ORDER_ID можно выполнить двумя способами:

  • собрать данные со всех узлов на один (привести к single-распределению);

  • собрать локально на узлах, если данные коллоцированы.

Таблица ORDERS уже распределена по аффинити-ключу ORDER_ID, но у таблицы ORDER_ITEMS другое распределение. Оператор IgniteExchange(distribution=[affinity[identity=RendezvousAffinityFunction, cacheId=xxxxxxxxxx][0]]) означает, что входящие в него данные нужно переслать узлам кластера с помощью аффинити-функции RendezvousAffinityFunction и распределения партиций по узлам кластера для кеша cacheId=xxxxxxxxxx. В качестве ключа в аффинити-функцию нужно передать колонку с нулевым индексом из входящих данных ([0] в конце distribution). Колонка с нулевым индексом будет соответствовать колонке с индексом 1 из таблицы, так как из таблицы выбираются всего три колонки начиная с индекса 1 (колонка с индексом 0 пропускается): requiredColumns=[{1, 2, 3}]. Соединение на каждом узле кластера выполнится коллоцировано и локально, результат соединения также будет распределен по узлам кластера по колонке ORDER_ID. Чтобы передать результат на узел-инициатор, потребуется еще один обмен с single-распределением.

IgniteTrimExchange#

Приведение broadcast-распределения данных в кластере к другому с помощью отбрасывания из выборки данных, которые не соответствуют новому распределению. IgniteTrimExchange используется для:

  • приведения к hash-/affinity-распределению;

  • коллокации данных;

  • в случае соединения broadcast-распределенных и hash-/affinity-распределенных таблиц.

Приведение к single-распределению не требует отдельного оператора. single-распределение можно получить из broadcast-распределения, выполнив запрос на одном узле, а не на всех. Приведение к random-распределению не имеет смысла, так как соединение таблиц на нескольких узлах кластера, у одной из которых random-распределение, невозможно.

Атрибуты:

  • distribution — распределение, к которому требуется привести.

SELECT * FROM items WHERE item_id in (SELECT item_id FROM bestsellers)
IgniteExchange(distribution=[single]): rowcount = 37500.0, cumulative cost = IgniteCost [rowCount=2575002.0, cpu=5450002.0, memory=6.0, io=2.0, network=300002.0], id = xxxx
  IgniteProject(ITEM_ID=[$0], NAME=[$1]): rowcount = 37500.0, cumulative cost = IgniteCost [rowCount=2537502.0, cpu=5412502.0, memory=6.0, io=2.0, network=2.0], id = xxxx
    IgniteMergeJoin(condition=[=($0, $2)], joinType=[inner], variablesSet=[[]], leftCollation=[[0 ASC-nulls-first]], rightCollation=[[0 ASC-nulls-first]]): rowcount = 37500.0, cumulative cost = IgniteCost [rowCount=2500002.0, cpu=5375002.0, memory=6.0, io=2.0, network=2.0], id = xxxx
      IgniteIndexScan(table=[[PUBLIC, ITEMS]], index=[ITEMS_ITEM_ID_IDX], inlineScan=[false], collation=[[0 ASC-nulls-first]]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500001.0, cpu=500001.0, memory=1.0, io=1.0, network=1.0], id = xxx
      IgniteTrimExchange(distribution=[affinity[identity=org.apache.ignite.cache.affinity.rendezvous.RendezvousAffinityFunction, cacheId=xxxxxxxxxx][0]]): rowcount = 250000.0, cumulative cost = IgniteCost [rowCount=1250001.0, cpu=1875001.0, memory=5.0, io=1.0, network=1.0], id = xxxx
        IgniteColocatedSortAggregate(group=[{0}], collation=[[0 ASC-nulls-first]]): rowcount = 250000.0, cumulative cost = IgniteCost [rowCount=1000001.0, cpu=875001.0, memory=5.0, io=1.0, network=1.0], id = xxxx
          IgniteIndexScan(table=[[PUBLIC, BESTSELLERS]], index=[BESTSELLERS_ITEM_ID_IDX], requiredColumns=[{0}], inlineScan=[true], collation=[[0 ASC-nulls-first]]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500001.0, cpu=375001.0, memory=1.0, io=1.0, network=1.0], id = xx

Пояснения к плану: у данных в таблице BESTSELLERS broadcast-распределение, они находятся на всех узлах кластера. IgniteTrimExchange приводит их к affinity-распределению по колонке ITEM_ID, чтобы выполнить коллоцированное соединение с таблицей ITEMS. У результата соединения также будет affinity-распределение по колонке ITEM_ID.

Агрегаты#

IgniteMapHashAggregate, IgniteReduceHashAggregate#

Распределенные агрегаты на основе хеш-таблицы. Состоят из двух фаз:

  • Map выполняется на всех узлах кластера;

  • Reduce выполняется на узле-инициаторе.

Между фазами Map и Reduce всегда присутствует IgniteExchange с single-распределением.

Атрибуты:

  • group — колонки группировки;

  • список агрегатов для вычисления.

SELECT order_date, AVG(total_sum), MIN(total_sum) FROM orders GROUP BY order_date
IgniteReduceHashAggregate(rowType=[RecordType(DATE ORDER_DATE, DECIMAL(32767, 0) EXPR$1, DECIMAL(32767, 0) EXPR$2)], group=[{0}], EXPR$1=[AVG($1)], EXPR$2=[MIN($1)]): rowcount = 250000.0, cumulative cost = IgniteCost [rowCount=1500000.0, cpu=1500000.0, memory=3500010.0, io=0.0, network=3000000.0], id = xx
  IgniteExchange(distribution=[single]): rowcount = 250000.0, cumulative cost = IgniteCost [rowCount=1250000.0, cpu=1250000.0, memory=3500000.0, io=0.0, network=3000000.0], id = xx
    IgniteMapHashAggregate(group=[{0}], EXPR$1=[AVG($1)], EXPR$2=[MIN($1)]): rowcount = 250000.0, cumulative cost = IgniteCost [rowCount=1000000.0, cpu=1000000.0, memory=3500000.0, io=0.0, network=0.0], id = xx
      IgniteTableScan(table=[[PUBLIC, ORDERS]], requiredColumns=[{1, 2}]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500000.0, cpu=500000.0, memory=0.0, io=0.0, network=0.0], id = xx

Пояснения к плану: на фазе Map происходит агрегирование локальных для узла данных и расчет предварительного результата. Например, если происходит группировка входящих данных в памяти по колонке с индексом 0:

  • для AVG сохраняется сумма и количество для каждого уникального значения колонки 0;

  • для MIN сохраняется минимальное значение для каждого уникального значения.

На фазе Reduce также в памяти с группировкой обрабатываются сгруппированные значения со всех узлов кластера. После того, как все данные получены, выдаются сгруппированные строки с постобработкой накопленных значений (например, AVG = сумма / количество).

IgniteColocatedHashAggregate#

Локальный агрегат на основе хеш-таблицы. Используется, если входящее распределение является:

  • single — все входящие данные доступны на узле-инициаторе, агрегация выполняется на узле-инициаторе;

  • broadcast — все входящие данные доступны на всех узлах кластера, агрегация выполняется одновременно на всех узлах кластера.

Атрибуты:

  • group — колонки группировки;

  • список агрегатов для вычисления.

Пример запроса, в котором данные обладают single-распределением уже после первой агрегации:

SELECT AVG(avg_sum) AS avg_per_day FROM (SELECT order_date, AVG(total_sum) avg_sum FROM orders GROUP BY order_date)
IgniteColocatedHashAggregate(group=[{}], AVG_PER_DAY=[AVG($0)]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=2000000.0, cpu=2000000.0, memory=2250010.0, io=0.0, network=3000000.0], id = xxx
  IgniteProject(AVG_SUM=[$1]): rowcount = 250000.0, cumulative cost = IgniteCost [rowCount=1750000.0, cpu=1750000.0, memory=2250005.0, io=0.0, network=3000000.0], id = xxx
    IgniteReduceHashAggregate(rowType=[RecordType(DATE ORDER_DATE, DECIMAL(32767, 0) AVG_SUM)], group=[{0}], AVG_SUM=[AVG($1)]): rowcount = 250000.0, cumulative cost = IgniteCost [rowCount=1500000.0, cpu=1500000.0, memory=2250005.0, io=0.0, network=3000000.0], id = xxx
      IgniteExchange(distribution=[single]): rowcount = 250000.0, cumulative cost = IgniteCost [rowCount=1250000.0, cpu=1250000.0, memory=2250000.0, io=0.0, network=3000000.0], id = xxx
        IgniteMapHashAggregate(group=[{0}], AVG_SUM=[AVG($1)]): rowcount = 250000.0, cumulative cost = IgniteCost [rowCount=1000000.0, cpu=1000000.0, memory=2250000.0, io=0.0, network=0.0], id = xxx
          IgniteTableScan(table=[[PUBLIC, ORDERS]], requiredColumns=[{1, 2}]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500000.0, cpu=500000.0, memory=0.0, io=0.0, network=0.0], id = xx

Пояснения к плану: в этом запросе данные обладают single-распределением уже после первой агрегации, поэтому вторую агрегацию можно выполнить с помощью IgniteColocatedHashAggregate только на узле-инициаторе запроса.

Пример запроса, в котором данные таблицы BESTSELLERS доступны на всех узлах кластера, но распределенная обработка для них не требуется:

SELECT MIN(item_id), MAX(item_id) FROM bestsellers
IgniteColocatedHashAggregate(group=[{}], EXPR$0=[MIN($0)], EXPR$1=[MAX($0)]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=1000001.0, cpu=875001.0, memory=11.0, io=1.0, network=1.0], id = xx
  IgniteIndexScan(table=[[PUBLIC, BESTSELLERS]], index=[BESTSELLERS_ITEM_ID_IDX], requiredColumns=[{0}], inlineScan=[true], collation=[[0 ASC-nulls-first]]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500001.0, cpu=375001.0, memory=1.0, io=1.0, network=1.0], id = xx

Пояснения к плану: в этом запросе данные таблицы BESTSELLERS доступны на всех узлах кластера, но распределенная обработка для них не требуется, поэтому запрос может выполниться на узле-инициаторе.

Пример запроса, в котором данные таблицы BESTSELLERS доступны на всех узлах кластера, агрегация также выполняется на всех узлах кластера и везде получается идентичный результат с broadcast-распределением:

SELECT * FROM items WHERE item_id = (SELECT MIN(item_id) FROM bestsellers WHERE description IS NOT NULL)
IgniteExchange(distribution=[single]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=1000005.0, cpu=2525007.000372529, memory=9.0, io=0.0, network=8.0], id = xxxx
  IgniteProject(ITEM_ID=[$0], NAME=[$1]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=1000004.0, cpu=2525006.000372529, memory=9.0, io=0.0, network=0.0], id = xxxx
    IgniteNestedLoopJoin(condition=[=($0, $2)], joinType=[inner], variablesSet=[[]]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=1000003.0, cpu=2525005.000372529, memory=9.0, io=0.0, network=0.0], id = xxxx
      IgniteTableScan(table=[[PUBLIC, ITEMS]]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500000.0, cpu=500000.0, memory=0.0, io=0.0, network=0.0], id = xxx
      IgniteTrimExchange(distribution=[affinity[identity=org.apache.ignite.cache.affinity.rendezvous.RendezvousAffinityFunction, cacheId=-1004180878][0]]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=3.0, cpu=25005.00037252903, memory=5.0, io=0.0, network=0.0], id = xxxx
        IgniteColocatedHashAggregate(group=[{}], EXPR$0=[MIN($0)]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=2.0, cpu=25001.00037252903, memory=5.0, io=0.0, network=0.0], id = xxxx
          IgniteIndexBound(table=[[PUBLIC, BESTSELLERS]], index=[BESTSELLERS_ITEM_ID_IDX], first=[true], requiredCols=[{0}]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=1.0, cpu=25000.00037252903, memory=0.0, io=0.0, network=0.0], id = xxx

Пояснения к плану: в этом запросе данные таблицы BESTSELLERS доступны на всех узлах кластера, агрегация (вычисление MIN(ITEM_ID)) также выполняется на всех узлах кластера и везде получается идентичный результат с broadcast-распределением. Далее результат остается только на тех узлах, которые соответствуют аффинити-функции для таблицы ITEMS, и происходит соединение с таблицей.

IgniteMapSortAggregate, IgniteReduceSortAggregate#

Распределенные агрегаты по отсортированным данным. Используются, если входящая сортировка соответствует группируемым колонкам (то есть агрегацию можно выполнить поточно, не накапливая строки в памяти). Как и соответствующий hash-агрегат, IgniteMapSortAggregate и IgniteReduceSortAggregate состоят из двух фаз:

  • Map выполняется на всех узлах кластера;

  • Reduce выполняется на узле-инициаторе.

Между фазами Map и Reduce всегда присутствует IgniteExchange с single-распределением.

Атрибуты:

  • group — колонки группировки;

  • список агрегатов для вычисления;

  • collation — порядок данных.

SELECT SUM(amount) FROM order_items GROUP BY order_id
IgniteProject(EXPR$0=[$1]): rowcount = 250000.0, cumulative cost = IgniteCost [rowCount=1750001.0, cpu=1750001.0, memory=10.0, io=1.0, network=2000001.0], id = xx
  IgniteReduceSortAggregate(rowType=[RecordType(INTEGER ORDER_ID, BIGINT EXPR$0)], group=[{0}], EXPR$0=[SUM($1)], collation=[[0 ASC-nulls-first]]): rowcount = 250000.0, cumulative cost = IgniteCost [rowCount=1500001.0, cpu=1500001.0, memory=10.0, io=1.0, network=2000001.0], id = xx
    IgniteExchange(distribution=[single]): rowcount = 250000.0, cumulative cost = IgniteCost [rowCount=1250001.0, cpu=1250001.0, memory=10.0, io=1.0, network=2000001.0], id = xx
      IgniteMapSortAggregate(group=[{0}], EXPR$0=[SUM($1)], collation=[[0 ASC-nulls-first]]): rowcount = 250000.0, cumulative cost = IgniteCost [rowCount=1000001.0, cpu=1000001.0, memory=10.0, io=1.0, network=1.0], id = xx
        IgniteIndexScan(table=[[PUBLIC, ORDER_ITEMS]], index=[ORDER_ITEMS_ORDER_ID_IDX], requiredColumns=[{1, 3}], inlineScan=[false], collation=[[1 ASC-nulls-first]]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500001.0, cpu=500001.0, memory=1.0, io=1.0, network=1.0], id = xx

Пояснения к плану: в таблице ORDER_ITEMS есть индекс по полю ORDER_ID, по которому нужна группировка. На фазе Map происходит суммирование значений AMOUNT для идентичных ORDER_ID. Как только приходит новый ORDER_ID, старое значение с суммой поднимается наверх (пересылается на узел-инициатор). На фазе Reduce таким же образом обрабатываются отсортированные данные из нескольких потоков (нескольких узлов кластера).

IgniteColocatedSortAggregate#

Локальный агрегат по отсортированным данным. Используется, если входящая сортировка соответствует группируемым колонкам (то есть агрегацию можно выполнить поточно, не накапливая строки в памяти). Как и соответствующий hash-агрегат, IgniteColocatedSortAggregate используется, если входящее распределение является:

  • single — все входящие данные доступны на узле-инициаторе, агрегация выполняется на узле-инициаторе;

  • broadcast — все входящие данные доступны на всех узлах кластера, агрегация выполняется одновременно на всех узлах кластера.

Атрибуты:

  • group — колонки группировки;

  • список агрегатов для вычисления;

  • collation — порядок данных.

SELECT * FROM items WHERE item_id in (SELECT item_id FROM bestsellers)
IgniteExchange(distribution=[single]): rowcount = 37500.0, cumulative cost = IgniteCost [rowCount=2575002.0, cpu=5450002.0, memory=6.0, io=2.0, network=300002.0], id = xxxx
  IgniteProject(ITEM_ID=[$0], NAME=[$1]): rowcount = 37500.0, cumulative cost = IgniteCost [rowCount=2537502.0, cpu=5412502.0, memory=6.0, io=2.0, network=2.0], id = xxxx
    IgniteMergeJoin(condition=[=($0, $2)], joinType=[inner], variablesSet=[[]], leftCollation=[[0 ASC-nulls-first]], rightCollation=[[0 ASC-nulls-first]]): rowcount = 37500.0, cumulative cost = IgniteCost [rowCount=2500002.0, cpu=5375002.0, memory=6.0, io=2.0, network=2.0], id = xxxx
      IgniteIndexScan(table=[[PUBLIC, ITEMS]], index=[ITEMS_ITEM_ID_IDX], inlineScan=[false], collation=[[0 ASC-nulls-first]]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500001.0, cpu=500001.0, memory=1.0, io=1.0, network=1.0], id = xxx
      IgniteTrimExchange(distribution=[affinity[identity=org.apache.ignite.cache.affinity.rendezvous.RendezvousAffinityFunction, cacheId=-1004180878][0]]): rowcount = 250000.0, cumulative cost = IgniteCost [rowCount=1250001.0, cpu=1875001.0, memory=5.0, io=1.0, network=1.0], id = xxxx
        IgniteColocatedSortAggregate(group=[{0}], collation=[[0 ASC-nulls-first]]): rowcount = 250000.0, cumulative cost = IgniteCost [rowCount=1000001.0, cpu=875001.0, memory=5.0, io=1.0, network=1.0], id = xxxx
          IgniteIndexScan(table=[[PUBLIC, BESTSELLERS]], index=[BESTSELLERS_ITEM_ID_IDX], requiredColumns=[{0}], inlineScan=[true], collation=[[0 ASC-nulls-first]]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500001.0, cpu=375001.0, memory=1.0, io=1.0, network=1.0], id = xx

Пояснения к плану: хотя в данном запросе нет явного агрегата, подзапрос для оператора IN требует уникальных значений item_id во входящей выборке. Эта уникальность достигается с помощью оператора DISTINCT, который также выполняется как агрегация, без агрегирующей функции с группировкой по требующим уникальности полям.

Операции над множествами#

IgniteMapMinus, IgniteReduceMinus, IgniteMapIntersect, IgniteReduceIntersect#

Распределенные операции разности (EXCEPT) и пересечения (INTERSECT) множеств на основе хеш-таблицы. Состоят из двух фаз:

  • Map выполняется на всех узлах кластера;

  • Reduce выполняется на узле-инициаторе.

Между фазами Map и Reduce всегда присутствует IgniteExchange с single-распределением. IgniteMapMinus, IgniteReduceMinus, IgniteMapIntersect и IgniteReduceIntersect работают по аналогии с хеш-агрегатами. У Map-операторов может быть два и более дочерних элемента.

Атрибуты:

  • all — для операций EXCEPT ALL, INTERSECT ALL.

SELECT order_id FROM orders WHERE order_date = ? EXCEPT SELECT order_id FROM order_items
IgniteReduceIntersect(all=[false], rowType=[RecordType(INTEGER $f0)]): rowcount = 37500.0, cumulative cost = IgniteCost [rowCount=1650001.0, cpu=8475001.0, memory=3675001.0, io=1.0, network=300001.0], id = xxx
  IgniteExchange(distribution=[single]): rowcount = 37500.0, cumulative cost = IgniteCost [rowCount=1612501.0, cpu=8100001.0, memory=3450001.0, io=1.0, network=300001.0], id = xxx
    IgniteMapIntersect(all=[false]): rowcount = 37500.0, cumulative cost = IgniteCost [rowCount=1575001.0, cpu=8062501.0, memory=3450001.0, io=1.0, network=1.0], id = xxx
      IgniteTableScan(table=[[PUBLIC, ORDERS]], filters=[=($t1, ?0)], projects=[[$t0]], requiredColumns=[{0, 1}]): rowcount = 75000.0, cumulative cost = IgniteCost [rowCount=500000.0, cpu=2000000.0, memory=0.0, io=0.0, network=0.0], id = xx
      IgniteIndexScan(table=[[PUBLIC, ORDER_ITEMS]], index=[ORDER_ITEMS_ORDER_ID_IDX], requiredColumns=[{1}], inlineScan=[true], collation=[[1 ASC-nulls-first]]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500001.0, cpu=312501.0, memory=1.0, io=1.0, network=1.0], id = xx

IgniteColocatedMinus, IgniteColocatedIntersect#

Локальные операции разности (EXCEPT) и пересечения (INTERSECT) множеств на основе хеш-таблицы. Используются, если входящее распределение является:

  • single — все входящие данные доступны на узле-инициаторе, агрегация выполняется на узле-инициаторе;

  • broadcast — все входящие данные доступны на всех узлах кластера, агрегация выполняется одновременно на всех узлах кластера.

Атрибуты:

  • all — для операций EXCEPT ALL, INTERSECT ALL.

SELECT item_id FROM items EXCEPT SELECT item_id FROM bestsellers
IgniteColocatedMinus(all=[false]): rowcount = 250000.0, cumulative cost = IgniteCost [rowCount=2500002.0, cpu=1.1250002E7, memory=6000002.0, io=2.0, network=2000002.0], id = xxx
  IgniteExchange(distribution=[single]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=1000001.0, cpu=875001.0, memory=1.0, io=1.0, network=2000001.0], id = xxx
    IgniteIndexScan(table=[[PUBLIC, ITEMS]], index=[ITEMS_ITEM_ID_IDX], requiredColumns=[{0}], inlineScan=[true], collation=[[0 ASC-nulls-first]]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500001.0, cpu=375001.0, memory=1.0, io=1.0, network=1.0], id = xx
  IgniteIndexScan(table=[[PUBLIC, BESTSELLERS]], index=[BESTSELLERS_ITEM_ID_IDX], requiredColumns=[{0}], inlineScan=[true], collation=[[0 ASC-nulls-first]]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500001.0, cpu=375001.0, memory=1.0, io=1.0, network=1.0], id = xx

IgniteUnionAll#

Объединение множеств. Внутри оператора агрегация не выполняется, но если требуется объединение множеств с контролем уникальности строк (оператор UNION без ключевого слова ALL), поверх оператора IgniteUnionAll добавляется агрегация (DISTINCT).

Атрибуты:

  • all — для операции UNION ALL (всегда установлено значение true).

SELECT item_id FROM order_items UNION SELECT item_id FROM bestsellers
IgniteColocatedHashAggregate(group=[{0}]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=3500001.0, cpu=3375001.0, memory=2000001.0, io=1.0, network=2000001.0], id = xxx
  IgniteUnionAll(all=[true]): rowcount = 1000000.0, cumulative cost = IgniteCost [rowCount=2500001.0, cpu=2375001.0, memory=1.0, io=1.0, network=2000001.0], id = xxx
    IgniteExchange(distribution=[single]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=1000000.0, cpu=1000000.0, memory=0.0, io=0.0, network=2000000.0], id = xxx
      IgniteTableScan(table=[[PUBLIC, ORDER_ITEMS]], requiredColumns=[{2}]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500000.0, cpu=500000.0, memory=0.0, io=0.0, network=0.0], id = xx
    IgniteIndexScan(table=[[PUBLIC, BESTSELLERS]], index=[BESTSELLERS_ITEM_ID_IDX], requiredColumns=[{0}], inlineScan=[true], collation=[[0 ASC-nulls-first]]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500001.0, cpu=375001.0, memory=1.0, io=1.0, network=1.0], id = xx

Пояснения к плану: так как требуется уникальность строк, поверх оператора IgniteUnionAll добавлена агрегация.

Соединения таблиц (join)#

У соединений всегда есть два дочерних узла дерева. Для удобства их называют левое плечо (первый дочерний узел) и правое плечо (второй дочерний узел). Ключ соединения — набор колонок из левого и правого плеч, которые сравниваются с помощью оператора «равенство» и объединены конъюнкцией в условиях соединения. Соединения возможны только для коллоцированных таблиц, если ключом соединения является аффинити-колонка. Если таблицы не коллоцированы, соединение выполняется с помощью узла-инициатора.

IgniteNestedLoopJoin#

Соединение с использованием вложенных циклов. Правое плечо оператора материализуется (все пришедшие строки сохраняются в памяти), затем для каждой строки из левого плеча проверяются на соответствие условию соединения все строки из правого плеча (то есть обрабатывается прямое (декартово) произведение левого и правого плеч).

Атрибуты:

  • condition — условие соединения;

  • joinType — тип соединения (INNER, LEFT, RIGHT, FULL, ANTI, SEMI).

SELECT name FROM items i LEFT OUTER JOIN bestsellers b ON (i.item_id = b.item_id OR i.name like '%best%')
IIgniteProject(NAME=[$1]): rowcount = 6.25E10, cumulative cost = IgniteCost [rowCount=3.12501500001E11, cpu=1.062501375001E12, memory=2000001.0, io=1.0, network=4000001.0], id = xxx
  IgniteNestedLoopJoin(condition=[OR(=($0, $2), LIKE($1, _UTF-8'%best%'))], joinType=[left], variablesSet=[[]]): rowcount = 6.25E10, cumulative cost = IgniteCost [rowCount=2.50001500001E11, cpu=1.000001375001E12, memory=2000001.0, io=1.0, network=4000001.0], id = xxx
    IgniteExchange(distribution=[single]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=1000000.0, cpu=1000000.0, memory=0.0, io=0.0, network=4000000.0], id = xxx
      IgniteTableScan(table=[[PUBLIC, ITEMS]]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500000.0, cpu=500000.0, memory=0.0, io=0.0, network=0.0], id = xx
    IgniteIndexScan(table=[[PUBLIC, BESTSELLERS]], index=[BESTSELLERS_ITEM_ID_IDX], requiredColumns=[{0}], inlineScan=[true], collation=[[0 ASC-nulls-first]]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500001.0, cpu=375001.0, memory=1.0, io=1.0, network=1.0], id = xx

Пояснения к плану: в этом запросе нет корректного ключа (условие содержит OR), по которому было бы можно использовать более оптимальный тип соединения, поэтому единственный вариант — перебор полного (декартова) произведения двух входящих потоков данных.

IgniteMergeJoin#

Соединение слиянием. Возможно, когда оба плеча отсортированы по ключу соединения.

Атрибуты:

  • condition — условие соединения;

  • joinType — тип соединения (INNER, LEFT, RIGHT, FULL, ANTI, SEMI).

SELECT o.order_date, oi.item_id FROM orders o JOIN order_items oi ON (o.order_id = oi.order_id) WHERE oi.amount = 1
IgniteExchange(distribution=[single]): rowcount = 11250.0, cumulative cost = IgniteCost [rowCount=1672502.0, cpu=4897502.0, memory=2.0, io=2.0, network=990002.0], id = xxxxx
  IgniteProject(ORDER_DATE=[$1], ITEM_ID=[$3]): rowcount = 11250.0, cumulative cost = IgniteCost [rowCount=1661252.0, cpu=4886252.0, memory=2.0, io=2.0, network=900002.0], id = 13886
    IgniteMergeJoin(condition=[=($0, $2)], joinType=[inner], variablesSet=[[]], leftCollation=[[0 ASC-nulls-first]], rightCollation=[[0 ASC-nulls-first]]): rowcount = 11250.0, cumulative cost = IgniteCost [rowCount=1650002.0, cpu=4875002.0, memory=2.0, io=2.0, network=900002.0], id = xxxxx
      IgniteIndexScan(table=[[PUBLIC, ORDERS]], index=[ORDER_ID_IDX], requiredColumns=[{0, 1}], inlineScan=[false], collation=[[0 ASC-nulls-first]]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500001.0, cpu=500001.0, memory=1.0, io=1.0, network=1.0], id = xxx
      IgniteExchange(distribution=[affinity[identity=org.apache.ignite.cache.affinity.rendezvous.RendezvousAffinityFunction, cacheId=-38276728][0]]): rowcount = 75000.0, cumulative cost = IgniteCost [rowCount=575001.0, cpu=2075001.0, memory=1.0, io=1.0, network=900001.0], id = xxxxx
        IgniteIndexScan(table=[[PUBLIC, ORDER_ITEMS]], index=[ORDER_ITEMS_ORDER_ID_IDX], filters=[=($t2, 1)], requiredColumns=[{1, 2, 3}], inlineScan=[false], collation=[[1 ASC-nulls-first]]): rowcount = 75000.0, cumulative cost = IgniteCost [rowCount=500001.0, cpu=2000001.0, memory=1.0, io=1.0, network=1.0], id = xxx

Пояснения к плану: в таблицах ORDERS и ORDER_ITEMS есть индексы по полю ORDER_ID, которое является ключом соединения, поэтому для данного запроса возможно соединение слиянием этих таблиц.

IgniteCorrelatedNestedLoopJoin#

Соединение с использованием корреляционной переменной. В данном типе соединения для каждой строки левого плеча перематывается и повторно итерируется правое плечо. Строка левого плеча устанавливается в качестве корреляционной переменной, а в правом плече выполняется соответствующая обработка этой переменной. Данный тип соединения возможен, когда у правого плеча есть возможность перемотки (характеристика rewindability = true).

Атрибуты:

  • condition — условие соединения;

  • joinType — тип соединения (INNER, LEFT, RIGHT, FULL, ANTI, SEMI);

  • correlationVariables — набор корреляционных переменных.

SELECT (SELECT description FROM bestsellers b WHERE i.item_id = b.item_id) FROM items i WHERE name like '%best%'
IgniteProject(EXPR$0=[$2]): rowcount = 125000.0, cumulative cost = IgniteCost [rowCount=1.8751E10, cpu=4.6882795886266525E10, memory=750000.0, io=125000.0, network=1125000.0], id = xxx
  IgniteCorrelatedNestedLoopJoin(condition=[true], joinType=[left], variablesSet=[[$cor0]], variablesSet=[[0]], correlationVariables=[[$cor0]]): rowcount = 125000.0, cumulative cost = IgniteCost [rowCount=1.8750875E10, cpu=4.6882670886266525E10, memory=750000.0, io=125000.0, network=1125000.0], id = xxx
    IgniteExchange(distribution=[single]): rowcount = 125000.0, cumulative cost = IgniteCost [rowCount=625000.0, cpu=2125000.0, memory=0.0, io=0.0, network=1000000.0], id = xxx
      IgniteTableScan(table=[[PUBLIC, ITEMS]], filters=[LIKE($t1, _UTF-8'%best%')]): rowcount = 125000.0, cumulative cost = IgniteCost [rowCount=500000.0, cpu=2000000.0, memory=0.0, io=0.0, network=0.0], id = xx
    IgniteColocatedHashAggregate(group=[{}], agg#0=[SINGLE_VALUE($0)]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=150001.0, cpu=375040.3670901322, memory=6.0, io=1.0, network=1.0], id = xxx
      IgniteIndexScan(table=[[PUBLIC, BESTSELLERS]], index=[BESTSELLERS_ITEM_ID_IDX], filters=[=($cor0.ITEM_ID, $t0)], projects=[[$t1]], requiredColumns=[{0, 1}], searchBounds=[[ExactBounds [bound=$cor0.ITEM_ID], null]], inlineScan=[false], collation=[[0 ASC-nulls-first]]): rowcount = 75000.0, cumulative cost = IgniteCost [rowCount=75001.0, cpu=300040.3670901322, memory=1.0, io=1.0, network=1.0], id = xxx

Пояснения к плану: соединение устанавливает корреляционную переменную $cor0 для каждой отфильтрованной строки из таблицы ITEMS, затем колонка ITEM_ID из этой строки используется в правом плече для поиска по индексу BESTSELLERS_ITEM_ID_IDX (устанавливается граница индекса bound=$cor0.ITEM_ID).

Spool-операторы#

Spool-операторы предназначены для кеширования полного набора данных, которые получены из нижележащего источника, чтобы обеспечить возможность перемотки этих данных. Spool-операторы добавляются после IgniteExchange (так как этот оператор убирает свойство «перематываемость») и перед правым плечом IgniteCorrelatedNestedLoopJoin (так как этот оператор требует свойство «перематываемость» на правом плече).

IgniteTableSpool#

Spool-оператор без возможности быстрой фильтрации (хранилище набора строк в памяти).

Атрибуты:

  • readTypeLAZY (как только данные появляются, они передаются вышестоящему узлу), EAGER (данные не передаются вышестоящему узлу, пока не получено все от нижестоящего узла; в некоторых случаях используется для оператора модификации);

  • writeType — всегда EAGER (не используется в DataGrid).

SELECT (SELECT name FROM items i WHERE i.name like b.description) FROM bestsellers b
IgniteProject(EXPR$0=[$2]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=1.0625015E12, cpu=1.812503E12, memory=1.0000025E12, io=0.0, network=1.0E12], id = xxx
  IgniteCorrelatedNestedLoopJoin(condition=[true], joinType=[left], variablesSet=[[$cor0]], variablesSet=[[0]], correlationVariables=[[$cor0]]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=1.062501E12, cpu=1.8125025E12, memory=1.0000025E12, io=0.0, network=1.0E12], id = xxx
    IgniteTableScan(table=[[PUBLIC, BESTSELLERS]]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500000.0, cpu=500000.0, memory=0.0, io=0.0, network=0.0], id = xx
    IgniteColocatedHashAggregate(group=[{}], agg#0=[SINGLE_VALUE($0)]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=2125000.0, cpu=3625000.0, memory=2000005.0, io=0.0, network=2000000.0], id = xxx
      IgniteFilter(condition=[LIKE($0, $cor0.DESCRIPTION)]): rowcount = 125000.0, cumulative cost = IgniteCost [rowCount=2000000.0, cpu=3500000.0, memory=2000000.0, io=0.0, network=2000000.0], id = xxx
        IgniteTableSpool(readType=[LAZY], writeType=[EAGER]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=1500000.0, cpu=1500000.0, memory=2000000.0, io=0.0, network=2000000.0], id = xxx
          IgniteExchange(distribution=[single]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=1000000.0, cpu=1000000.0, memory=0.0, io=0.0, network=2000000.0], id = xxx
            IgniteTableScan(table=[[PUBLIC, ITEMS]], requiredColumns=[{1}]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500000.0, cpu=500000.0, memory=0.0, io=0.0, network=0.0], id = xxx

Пояснения к плану: нет возможности быстрой фильтрации по условию i.name like b.description, поэтому используется обычный Spool-оператор с затратами на обработку O(n * m).

IgniteSortedIndexSpool#

Spool-оператор с отсортированными данными. Поддерживает фильтрацию по полю за время O(n), в том числе поиск по диапазону значений.

Атрибуты:

  • readType — всегда LAZY;

  • writeType — всегда EAGER (не используется в DataGrid);

  • condition — условие;

  • collation — порядок данных;

  • searchBounds — границы поиска.

SELECT (SELECT name FROM items i WHERE i.item_id > b.item_id) FROM bestsellers b
IgniteProject(EXPR$0=[$2]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=1.000002E12, cpu=7.50023183545066E11, memory=2.000003E12, io=500000.0, network=2.0000005E12], id = xxx
  IgniteCorrelatedNestedLoopJoin(condition=[true], joinType=[left], variablesSet=[[$cor0]], variablesSet=[[0]], correlationVariables=[[$cor0]]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=1.0000015E12, cpu=7.50022683545066E11, memory=2.000003E12, io=500000.0, network=2.0000005E12], id = xxx
    IgniteTableScan(table=[[PUBLIC, BESTSELLERS]]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500000.0, cpu=500000.0, memory=0.0, io=0.0, network=0.0], id = xx
    IgniteColocatedHashAggregate(group=[{}], agg#0=[SINGLE_VALUE($0)]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=2000001.0, cpu=1500040.367090132, memory=4000006.0, io=1.0, network=4000001.0], id = xxx
      IgniteProject(NAME=[$1]): rowcount = 250000.0, cumulative cost = IgniteCost [rowCount=1750001.0, cpu=1250040.367090132, memory=4000001.0, io=1.0, network=4000001.0], id = xxx
        IgniteSortedIndexSpool(readType=[LAZY], writeType=[EAGER], condition=[>($0, $cor0.ITEM_ID)], collation=[[0 ASC-nulls-first]], searchBounds=[[RangeBounds [lowerBound=$cor0.ITEM_ID, upperBound=null, lowerInclude=false, upperInclude=true], null]]): rowcount = 250000.0, cumulative cost = IgniteCost [rowCount=1500001.0, cpu=1000040.3670901322, memory=4000001.0, io=1.0, network=4000001.0], id = xxx
          IgniteExchange(distribution=[single]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=1000001.0, cpu=1000001.0, memory=1.0, io=1.0, network=4000001.0], id = xxx
            IgniteIndexScan(table=[[PUBLIC, ITEMS]], index=[ITEMS_ITEM_ID_IDX], inlineScan=[false], collation=[[0 ASC-nulls-first]]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500001.0, cpu=500001.0, memory=1.0, io=1.0, network=1.0], id = xxx

Пояснения к плану: фильтрация происходит по диапазону значений, поэтому можно использовать Spool-оператор с сортировкой.

IgniteHashIndexSpool#

Spool-оператор на основе хеш-таблицы. Поддерживает фильтрацию по полю, если задано условие «равенство».

Атрибуты:

  • readType — всегда LAZY;

  • writeType — всегда EAGER (не используется в DataGrid);

  • searchRow — ключ для поиска в хеш-таблице;

  • condition — условие (может быть более жестким, чем searchRow);

  • allowNulls — поддержка поиска с учетом NULL-значений (для оператора IS NOT DISTINCT FROM).

SELECT (SELECT name FROM items i WHERE i.item_id like b.item_id) FROM bestsellers b
IgniteProject(EXPR$0=[$2]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=1.2500015E12, cpu=1.000008E12, memory=2.0000025E12, io=0.0, network=2.0E12], id = xxx
  IgniteCorrelatedNestedLoopJoin(condition=[true], joinType=[left], variablesSet=[[$cor0]], variablesSet=[[0]], correlationVariables=[[$cor0]]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=1.250001E12, cpu=1.0000075E12, memory=2.0000025E12, io=0.0, network=2.0E12], id = xxx
    IgniteTableScan(table=[[PUBLIC, BESTSELLERS]]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500000.0, cpu=500000.0, memory=0.0, io=0.0, network=0.0], id = xx
    IgniteColocatedHashAggregate(group=[{}], agg#0=[SINGLE_VALUE($0)]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=2500000.0, cpu=2000010.0, memory=4000005.0, io=0.0, network=4000000.0], id = xxx
      IgniteProject(NAME=[$1]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=2000000.0, cpu=1500010.0, memory=4000000.0, io=0.0, network=4000000.0], id = xxx
        IgniteHashIndexSpool(readType=[LAZY], writeType=[EAGER], searchRow=[[$cor0.ITEM_ID, null]], condition=[=($0, $cor0.ITEM_ID)], allowNulls=[false]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=1500000.0, cpu=1000010.0, memory=4000000.0, io=0.0, network=4000000.0], id = xxx
          IgniteExchange(distribution=[single]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=1000000.0, cpu=1000000.0, memory=0.0, io=0.0, network=4000000.0], id = xxx
            IgniteTableScan(table=[[PUBLIC, ITEMS]]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500000.0, cpu=500000.0, memory=0.0, io=0.0, network=0.0], id = xxx

Пояснения к плану: фильтрация происходит по конкретному значению, поэтому можно использовать Spool-оператор с хеш-таблицей.

Разворачивание/сворачивание коллекций#

IgniteCollect#

Сворачивание входящего набора строк в коллекцию.

Атрибуты:

  • field — имя поля в результирующем типе;

  • collectionType — тип коллекции (ARRAY/MAP).

SELECT array(SELECT item_id FROM bestsellers)
IgniteProject(EXPR$0=[$1]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=2000002.0, cpu=3375002.0, memory=1.0, io=1.0, network=1.0], id = xxx
  IgniteCorrelatedNestedLoopJoin(condition=[true], joinType=[inner], variablesSet=[[$cor1]], variablesSet=[[1]], correlationVariables=[[$cor1]]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=1500002.0, cpu=2875002.0, memory=1.0, io=1.0, network=1.0], id = xxx
    IgniteValues(tuples=[[{ 0 }]]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=1.0, cpu=1.0, memory=0.0, io=0.0, network=0.0], id = xx
    IgniteCollect(field=[x], collectionType=[ARRAY]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=1000001.0, cpu=875001.0, memory=1.0, io=1.0, network=1.0], id = xxx
      IgniteIndexScan(table=[[PUBLIC, BESTSELLERS]], index=[BESTSELLERS_ITEM_ID_IDX], requiredColumns=[{0}], inlineScan=[true], collation=[[0 ASC-nulls-first]]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500001.0, cpu=375001.0, memory=1.0, io=1.0, network=1.0], id = xx

Пояснения к плану: так как подзапросы разворачиваются с помощью соединения таблиц (join), в данном плане естьjoin, но на конечный результат он не влияет.

IgniteUncollect#

Разворачивание коллекции в набор строк.

Атрибуты:

Отсутствуют.

SELECT * FROM UNNEST(ARRAY[1, 2, 3])
IgniteUncollect: rowcount = 1.0, cumulative cost = IgniteCost [rowCount=3.0, cpu=3.0, memory=0.0, io=0.0, network=0.0], id = xx
  IgniteProject(EXPR$0=[ARRAY(1, 2, 3)]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=2.0, cpu=2.0, memory=0.0, io=0.0, network=0.0], id = xx

Прочее#

IgniteFilter#

Фильтрация строк. Может сливаться с операторами IgniteTableScan и IgniteIndexScan для упрощения плана, поэтому условие по колонкам таблицы отдельного оператора IgniteFilter обычно не создает.

Атрибуты:

  • condition — условие.

SELECT * FROM (VALUES (0), (1)) AS t(a) WHERE a > 0
IgniteFilter(condition=[>($0, 0)]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=4.0, cpu=9.0, memory=0.0, io=0.0, network=0.0], id = xx
  IgniteValues(tuples=[[{ 0 }, { 1 }]]): rowcount = 2.0, cumulative cost = IgniteCost [rowCount=2.0, cpu=1.0, memory=0.0, io=0.0, network=0.0], id = xx

IgniteProject#

Проекция. Может сливаться с операторами IgniteTableScan и IgniteIndexScan для упрощения плана, поэтому проекция колонок таблицы отдельного оператора IgniteProject обычно не создает.

Атрибуты:

  • список выражений для колонок.

SELECT a + 1 FROM (VALUES (0), (1)) AS t(a)
IgniteProject(EXPR$0=[+($0, 1)]): rowcount = 2.0, cumulative cost = IgniteCost [rowCount=4.0, cpu=3.0, memory=0.0, io=0.0, network=0.0], id = xx
  IgniteValues(tuples=[[{ 0 }, { 1 }]]): rowcount = 2.0, cumulative cost = IgniteCost [rowCount=2.0, cpu=1.0, memory=0.0, io=0.0, network=0.0], id = xx

IgniteSort#

Сортировка.

Атрибуты:

  • sort* — выражения для сортировки;

  • dir* — направление сортировки.

SELECT * FROM items ORDER BY name, item_id
IgniteExchange(distribution=[single]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=1500000.0, cpu=2.118354506610649E7, memory=4000000.0, io=0.0, network=4000000.0], id = xxx
  IgniteSort(sort0=[$1], sort1=[$0], dir0=[ASC-nulls-first], dir1=[ASC-nulls-first]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=1000000.0, cpu=2.068354506610649E7, memory=4000000.0, io=0.0, network=0.0], id = xxx
    IgniteTableScan(table=[[PUBLIC, ITEMS]], requiredColumns=[{0, 1}]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500000.0, cpu=500000.0, memory=0.0, io=0.0, network=0.0], id = xx

IgniteLimit#

Ограничение количества строк и смещение.

Атрибуты:

  • offset — смещение;

  • fetch — ограничение количества строк.

SELECT * FROM items ORDER BY name LIMIT 1 OFFSET 1
IgniteLimit(offset=[1], fetch=[1]): rowcount = 1.0, cumulative cost = IgniteCost [rowCount=1000003.0, cpu=2500003.0, memory=16.0, io=0.0, network=16.0], id = xx
  IgniteExchange(distribution=[single]): rowcount = 2.0, cumulative cost = IgniteCost [rowCount=1000002.0, cpu=2500002.0, memory=16.0, io=0.0, network=16.0], id = xx
    IgniteSort(sort0=[$1], dir0=[ASC-nulls-first], offset=[1], fetch=[1]): rowcount = 2.0, cumulative cost = IgniteCost [rowCount=1000000.0, cpu=2500000.0, memory=16.0, io=0.0, network=0.0], id = xx
      IgniteTableScan(table=[[PUBLIC, ITEMS]], requiredColumns=[{0, 1}]): rowcount = 500000.0, cumulative cost = IgniteCost [rowCount=500000.0, cpu=500000.0, memory=0.0, io=0.0, network=0.0], id = xx

IgniteTableModify#

Модификация данных. Используется для DML-операторов INSERT/UPDATE/DELETE/MERGE. Всегда выполняется на узле-инициаторе. Физически изменение выполняется с помощью операции KV API cache.invoke по ключу, который получен от источника данных снизу. Изменение входящих в ключ колонок невозможно.

Атрибуты:

  • table — таблица;

  • operation — тип операции;

  • updateColumnList — список колонок для изменения;

  • sourceExpressionList — список значений для установки.

UPDATE items SET name = 'Best item' WHERE item_id = ?
IgniteTableModify(table=[[PUBLIC, ITEMS]], operation=[UPDATE], updateColumnList=[[NAME]], sourceExpressionList=[[_UTF-8'Best item']], flattened=[false]): rowcount = 75000.0, cumulative cost = IgniteCost [rowCount=150002.0, cpu=375041.3670901322, memory=2.0, io=2.0, network=900002.0], id = xx
  IgniteExchange(distribution=[single]): rowcount = 75000.0, cumulative cost = IgniteCost [rowCount=150001.0, cpu=375040.3670901322, memory=1.0, io=1.0, network=900001.0], id = xx
    IgniteIndexScan(table=[[PUBLIC, ITEMS]], index=[ITEMS_ITEM_ID_IDX], filters=[=($t0, ?0)], projects=[[$t0, $t1, _UTF-8'Best item']], requiredColumns=[{0, 1}], searchBounds=[[ExactBounds [bound=?0], null]], inlineScan=[false], collation=[[0 ASC-nulls-first]]): rowcount = 75000.0, cumulative cost = IgniteCost [rowCount=75001.0, cpu=300040.3670901322, memory=1.0, io=1.0, network=1.0], id = xx