Структура фактического плана запроса
На диаграмме левая колонка показывает структуру стадий и операторов; ниже — их семантика, каналы связи и нетривиальные топологии.
Стадии
Ниже представлен запрос и его графическое представление:
SELECT count(*) cnt, l_orderkey, sum(l_quantity)
FROM lineitem
WHERE l_commitdate >= date("1993-01-01") and l_commitdate < date("1993-02-01")
GROUP BY l_orderkey
HAVING sum(l_quantity) > 200
ORDER BY cnt DESC, l_orderkey
В примере три вычислительных (compute) стадии, пронумерованные как 0, 1 и 2, а также отдельная стадия у таблицы lineitem (стадии чтения не нумеруются). Любую стадию можно временно выделить на схеме; здесь выделена стадия 1.
Интерактивные элементы на встроенных изображениях недоступны — см. Расположение информации в плане запроса.
На схеме — структура выполнения запроса, полученная после его компиляции: физические операторы, сгруппированные в стадии выполнения.
Стадия — основная единица плана; задачи стадии выполняются параллельно, часто на разных узлах. Данные передаются от стадии к стадии снизу вверх: от хранилища к результату.
Стадии бывают двух типов:
- Стадии чтения (scan), которые выполняются в хранилище (колоночные таблицы или строковые таблицы) имеют тёмный фон.
- Вычислительные стадии (compute) отображаются на светлом фоне и нумеруются последовательно, начиная с
0.
Примечание
Нумерованный список ниже на странице идёт сверху вниз, на графике стадии и операторы расположены снизу вверх. Сопоставлять пункты со схемой удобно, начиная с нижней части плана.
Из плана мы можем понять следующее:
- Таблица
lineitemсканируется целиком, потому что условиеWHEREне использует первичный ключ. - Фильтр по полю
l_commitdateвыполняется непосредственно в хранилище (pushdown предиката), это позволяет на раннем этапе отбросить ненужные данные и не пересылать их далее. - В стадии
0выполняется предварительная агрегация с группировкой поl_orderkey(конструкцияGROUP BYв тексте запроса). - Данные перераспределяются (hash shuffle) в стадию
1; подробнее — в разделе Каналы связи. - В стадии
1находится финальная агрегация, фильтр обозначенный в запросе в частиHAVINGи локальная (для каждого узла) сортировка. - Стадия
2объединяет результат обработки на каждом из узлов в глобально упорядоченную последовательность строк и возвращает клиенту в качестве результата работы запроса.
Совет
Неочевидный элемент схемы можно подсветить курсором: у большинства узлов есть всплывающая подсказка. Это не замена разделам документации, но удобно для быстрой ориентации и метрик.
Каналы связи
Канал связи на диаграмме показан пятиугольником с направлением потока. В этом примере выделен канал 0 и 1.
У линейного плана из четырёх стадий три соединения:
E— получение данных из хранилища (scan).H— hash shuffle, разбиение по одному или нескольким полям (список полей — во всплывающей подсказке по пятиугольнику).Me— слияние (merge) уже отсортированных потоков; для этого типа указываются поля сортировки.
Соединение состоит из двух частей:
- Выходной канал — исходящий трафик из стадии
0. Синяя стрелка вправо-вверх с номером стадии-получателя (здесь1). - Входной канал — приём в стадию
1. Зелёная стрелка влево-вверх с номером стадии-отправителя (0).
На графике успешно завершённого запроса объём данных на выходе канала и на входе должен совпадать: YDB гарантирует доставку между узлами без потерь и дубликатов.
Пока запрос выполняется, показатели входа и выхода могут расходиться из-за задержки буферов и асинхронного обновления статистики. На графике ошибочного запроса расхождения тоже возможны.
Входной канал сопоставляется оператору, который принимает данные. В простом случае у стадии один вход: поток из канала подаётся в первый по порядку обработки (нижний в списке) оператор. В сложных структурах выполнения запроса зачастую несколько входов; у оператора справа дублируется номер стадии-источника исходящего канала (здесь 0), он подсвечивается при выборе канала. Разбор на примере с JOIN — в разделе ниже.
Что можно понять для выбранного канала между стадиями 0 и 1:
- Начинается в стадии
0, направлен в стадию1. - Проходит через hash shuffle.
- Входит в стадию
1со стадии0. - Подаётся в оператор агрегации.
Стадии чтения (scans) и вычисления (compute) связаны аналогичными каналы связи, имеющими немного другой набор метрик, чем каналы, связывающие две вычислительных (compute) стадии (см. выше). Планировщик YDB по возможности размещает связанные задачи на одном и том же узле, чтобы сократить пересылаемый по сети объём данных и ускорить выполнения запроса.
Объединение стадий
Ниже — запрос с JOIN двух таблиц.
PRAGMA ydb.OptimizerHints = 'JoinType(nation region shuffle)';
SELECT n_name
FROM nation
JOIN region ON nation.n_regionkey == region.r_regionkey
WHERE r_name = "AMERICA"
Таблицы region и nation небольшие (5 и 25 строк), запрос быстрый, трафик небольшой, но структура плана показательна. Для примера задан Shuffle Join через подсказку оптимизатора. Без этой прагмы план будет другим: при малом объёме данных стоимостной оптимизатор выберет иной план (его можно сравнить, убрав PRAGMA).
Стадия 2 (на схеме изначально выделена) получает данные от двух стадий с номерами 0 и 1. Стадии расположены по вертикали, и направление движения данных — снизу вверх, однако связи между ними не обязательно последовательные. Если у стадии есть несколько предшественников, слева появляется отступ, а входные H-каналы обозначаются пятиугольником
Данные объединяются в физическом операторе InnerJoin (эквивалент JOIN в тексте запроса). Указаны ключи n_regionkey и r_regionkey; справа — номера стадий 0 и 1.
Красные круги с числами 1 и 5 — напоминание из раздела Агрегаты: число ненулевых метрик (здесь — по трафику в каналах) меньше числа задач стадии. Часть задач не получила данных, частая причина этого — перекос данных или избыточный параллелизм.
На этой схеме следует обратить внимание на два типичных случая, которые вызывают перекос по данным (data skew) и, таким образом, замедляют выполнение запроса, потому что те задачи, которые получили меньший объём трафика, завершатся быстрее и будут ждать другие, которым досталось больше данных:
1. Неравномерное распределение данных в хранилище
В столбце Tasks видно, что у стадий чтения указан параллелизм 1 (по одному шарду на таблицы region и nation). У связанных с ними вычислительных (compute) стадий с номерами 0 и 1 планировщиком создано по 3 задачи. Все прочитанные из каждой таблицы строки были сгруппированы в один пакет (batch) и отправлены из колоночного шарда в одну из трёх задач (случайную). Оставшимся задачам ничего не досталось, поэтому они не отчитались в статистике по этому каналу.
Когда в таблице данных больше, и пересылка использует более одного сообщения, то обычно данные распределяются равномерно между вычислительными задачами. Однако, если при хранении данных возникает перекос (одни шарды содержат больше данных чем другие), то похожая неравномерность также может проявляться и плохо влиять на производительность, поэтому она выделяется на схеме, чтобы обратить на это внимание.
2. Неравномерное распределение ключей
Даже если данные в хранилище распределены равномерно, они перераспределяются в процессе последующей обработки, что также может вызывать перекос, вызванный алгоритмическими особенностями реализации.
Соединение hash shuffle используется для корректной обработки больших объёмов данных несколькими задачами одновременно. Смысл этого действия заключается в том, чтобы разделить всё множество строк на несколько непересекающихся групп таким образом, чтобы строки с одинаковыми значениями определённых колонок оказались бы в одной группе. Тогда каждая задача может независимо от других обрабатывать свою группу и получить правильный результат.
Такие колонки ещё называют ключевыми, потому что они используются в качестве ключей для hash shuffle и последующих операторов.
В данном примере hash shuffle используется для того, чтобы сгруппировать необходимые строки в задаче, реализующей оператор JOIN. Он реализован как вычисление значения хеш-функции от всех ключевых колонок и далее остатка от деления этого значения на общее количество таких задач.
Но общее количество различных значений (кардинальность) хеш-функции не может превышать общего количества различных ключей в наших данных, поэтому даже если данных изначально много и все они распределены в хранилище равномерно, но количество ключей невелико и сильно отличается для разных значений ключей, у нас снова возможен перекос, который мы наблюдаем на этом плане.
После применения условия WHERE из таблицы region остаётся всего одна строка, которая естественным образом попадает только в одну из созданных задач, о чём сигнализирует красный круг с отметкой 1 на входе стадии 2 (мы наблюдаем низкую кардинальность по ключу r_regionkey).
Из таблицы nation выбираются все 25 строк. Однако для ключа n_regionkey у нас есть всего 5 разных значений. Поэтому из всех 8 задач стадии 2 данные из таблицы nation приходят только в 5 (остальным задачам ничего не досталось). Все 8 задач выполняют оператор JOIN, однако непустой результат будет только у одной из них — именно той, в которую попала строка из таблицы region. Эта задача вернёт 5 выбранных строк таблицы nation, связанных с данным регионом. Поэтому на выходе стадии 2 тоже появляется красная отметка 1.
Множественные выходы
Примечание
В примерах этого раздела используется набор данных TPC-H. Схема таблиц и связи между ними описана в спецификации TPC-H. Некоторые примеры являются запросами из бенчмарка с адаптацией под YDB, другие специально подготовлены для иллюстрации излагаемого материала. Исходный текст SQL всех примеров позволяет самостоятельно повторить все описанные примеры.
Следующий случай выполнения запроса: одна таблица читается один раз, а результат расходится к двум (или более) потребителям, выше потоки снова сходятся.
Такую структуру удобно показывать клонами стадии: один экземпляр — основной, остальные — сокращённое представление и отдельный выход. Пример — TPC-H Q17, адаптированный под синтаксис YDB:
SELECT SUM(l_extendedprice) / 7.0 AS avg_yearly
FROM lineitem
CROSS JOIN part
CROSS JOIN (
SELECT l_partkey, 0.2 * AVG(l_quantity) AS quantity_threshold
FROM lineitem
GROUP BY l_partkey
) AS threshold
WHERE part.p_partkey = lineitem.l_partkey
AND p_brand = 'Brand#35'
AND p_container = 'LG DRUM'
AND l_quantity < quantity_threshold
AND part.p_partkey = threshold.l_partkey
Таблица lineitem читается один раз в стадии 0 и используется в стадиях 1 и 3. На плане стадия 0 показана дважды:
- Основной экземпляр — полный набор метрик: входы, CPU, память.
- Клоны — тот же цвет, что у стадий чтения, и одна метрика: выход, ведущий в нужную нижестоящую стадию.
Выбор любого экземпляра стадии 0 подсвечивает все её копии; но выходы выбираются независимо друг от друга, чтобы можно было понять, куда идёт каждый из них.
Стоимостной оптимизатор может объединять и дедуплицировать не только чтение из хранилища, но и целые подструктуры — это хорошо видно, например, на TPC-H Q21.
Составные структуры
В примерах выше была одна структура выполнения. Если в одном вызове можно передать несколько выражений SQL, то сервер выполнит их по очереди, и на диаграмме появятся несколько несвязанных структур, каждая из которых представляет отдельное выполнение запроса.
SELECT count(*) FROM lineitem;
SELECT count(*) FROM orders;
Бывают случаи, когда одно SQL-выражение приводит к появлению сразу нескольких ветвей вычислений: стоимостной оптимизатор может вынести часть операций в отдельный прекомпьют — самостоятельную цепочку, создающую промежуточный результат для последующего использования в основном процессе выполнения запроса. Например, так устроена TPC-H Q15 в синтаксисе YDB:
$revenue0 = (
SELECT l_suppkey AS supplier_no, sum(l_extendedprice * (1 - l_discount)) AS total_revenue
FROM lineitem
WHERE l_shipdate >= date('1996-01-01') AND l_shipdate < date('1996-01-01') + interval('P90D')
GROUP BY l_suppkey
);
SELECT s_suppkey, s_name, s_address, s_phone, total_revenue
FROM supplier
CROSS JOIN $revenue0 AS revenu0
CROSS JOIN (
SELECT max(total_revenue) AS max_total_revenue
FROM $revenue0
) as max_revenue
WHERE s_suppkey = supplier_no AND total_revenue = max_total_revenue
ORDER BY s_suppkey;
Связь прекомпьюта с основной структурой выполнения становится очевидной при выделении имени прекомпьюта или его выходного канала с результатом: одновременно подсвечивается и вход оператора, который использует этот результат. Такой вход отмечен символом P. В данном примере это стадия 3 в верхней структуре; данные из прекомпьюта поступают в правую часть Map Join.
В тяжёлых запросах (в том числе из TPC-DS) прекомпьютов несколько; они могут выполняться последовательно или параллельно, результаты подмешиваются в следующие стадии или порождают новые прекомпьюты. Формат плана рассчитан на то, чтобы по схеме можно было восстановить реальную структуру выполнения и соответствие задач на кластере YDB.