Как читать Explain Analyze и ускорять долгие JOIN-ы в Greenplum?

Ускорение медленного JOIN-а в Greenplum редко начинается с индекса — здесь работает другая логика. Любая оптимизация запроса Greenplum начинается с чтения плана: EXPLAIN ANALYZE показывает не только оценку планировщика, но и то, как данные физически перемещаются между сегментами кластера. Это перемещение — Motion — чаще всего и есть причина долгого JOIN-а.

В статье разберём, чем EXPLAIN ANALYZE в MPP-системе отличается от обычного PostgreSQL, что такое Broadcast Motion, Redistribute Motion и Gather Motion, и как на реальном примере DWH-запроса ускорить JOIN в 12 раз, поправив один DISTRIBUTED BY.
Аудит Arenadata DB Администрирование Arenadata DB

Особенности Explain Analyze в MPP-архитектуре

Greenplum и Arenadata DB (arenadata db) — MPP-системы с архитектурой shared nothing: каждый сегмент хранит свою часть данных и выполняет запрос независимо, без общей памяти и диска. Координатор строит план, рассылает его сегментам, результаты собираются обратно. Один SQL-запрос физически выполняется десятками или сотнями параллельных процессов — по одному на сегмент.
Отсюда главное отличие от обычного PostgreSQL: план делится на слайсы (slice) — независимые части, выполняющиеся параллельно, а между слайсами данные передаются по внутренней сети — интерконнекту (interconnect). Узлы Motion в плане показывают, где происходит эта передача, — это первое, на что стоит смотреть при разборе долгого запроса.

Как читать узлы (Nodes)

План — это дерево узлов (plan node), которое читают снизу вверх: глубокие узлы выполняются первыми, результат поднимается к корню. У узла два блока метрик: cost — оценка планировщика в условных единицах; rows и width — сколько строк и какого размера вернёт узел по оценке; actual time, rows и loops — реальное время, число строк и число запусков узла.
Чаще всего встречаются Seq Scan (последовательное чтение таблицы), Hash Join (join через хэш-таблицу) и Nested Loop (вложенный цикл, приемлемый только для небольших наборов строк). Nested Loop на многомиллионных таблицах — почти всегда признак устаревшей статистики или неудачного distribution key. Главный сигнал проблемы — расхождение оценки и факта по rows: значит, статистика не отражает реальность. Лечится это VACUUM ANALYZE.

Motions: главные особенности производительности в Greenplum

Motion — узел, который перемещает данные между сегментами по интерконнекту и обычно определяет, будет JOIN быстрым или растянется на десятки минут. В Greenplum три типа Motion.

Broadcast Motion

Broadcast Motion рассылает полную копию таблицы на каждый сегмент кластера — при 96 сегментах после Broadcast Motion каждый из них хранит у себя всю таблицу целиком. Планировщик выбирает broadcast motion greenplum обычно для небольших справочников: дешевле разослать копию маленькой таблицы, чем перераспределять по сети большую.
Проблема — когда планировщик из-за устаревшей статистики считает broadcast выгодным для таблицы, которая на самом деле уже не маленькая.

Redistribute Motion

Redistribute Motion перераспределяет строки таблицы по сегментам заново, хэшируя ключ джойна. Если обе таблицы в JOIN распределены не по тому столбцу, планировщику приходится перетасовать по сети одну или обе таблицы. Это самый дорогой вид Motion при больших объёмах — именно он чаще всего виноват в долгих JOIN-ах в DWH.

Gather Motion

Gather Motion собирает результаты со всех сегментов в одну точку — обычно на координатор, для агрегации, сортировки или отправки клиенту. Он практически всегда есть в плане и редко бывает узким местом сам по себе.

Broadcast Motion vs Redistribute Motion: как данные летят между сегментами Greenplum

Оптимизация долгого JOIN-а на реальном примере

Разберём кейс из DWH: витрина считает агрегат по покупателям, JOIN-я таблицу фактов заказов с таблицей скоринга клиентов. Запрос выполнялся почти 42 секунды.

Анализ первоначального плана

EXPLAIN ANALYZE
SELECT fact.customer_id, sum(fact.amount)
FROM fact_orders fact
JOIN dim_customer_scoring scoring
  ON fact.customer_id = scoring.customer_id
GROUP BY fact.customer_id;
 
Gather Motion 96:1  (slice3; segments: 96)  (cost=0.00..892345.12 rows=1 width=48)
  (actual time=41823.441..41823.512 rows=1 loops=1)
  ->  HashAggregate  (cost=845210.30..845210.31 rows=1 width=48)
        (actual time=41756.203..41756.204 rows=1 loops=1)
        Group Key: fact.customer_id
        ->  Hash Join  (cost=312500.00..812340.00 rows=1310000 width=32)
              (actual time=8210.552..40982.117 rows=1298742 loops=1)
              Hash Cond: (fact.customer_id = scoring.customer_id)
              ->  Redistribute Motion 96:96  (slice2; segments: 96)
                    (cost=0.00..401200.00 rows=8400000 width=24)
                    (actual time=0.041..21432.884 rows=8391204 loops=1)
                    Hash Key: fact.customer_id
                    ->  Seq Scan on fact_orders fact
                          (cost=0.00..250340.00 rows=8400000 width=24)
                          (actual time=0.012..3452.201 rows=8391204 loops=1)
              ->  Hash  (cost=185300.00..185300.00 rows=1310000 width=16)
                    (actual time=8195.312..8195.312 rows=1298742 loops=1)
                    ->  Redistribute Motion 96:96  (slice1; segments: 96)
                          (cost=0.00..185300.00 rows=1310000 width=16)
                          (actual time=0.038..6821.442 rows=1298742 loops=1)
                          Hash Key: scoring.customer_id
                          ->  Seq Scan on dim_customer_scoring scoring
                                (cost=0.00..92100.00 rows=1310000 width=16)
                                (actual time=0.009..1103.221 rows=1298742 loops=1)
Planner: Pivotal Optimizer (GPORCA)
Planning time: 24.671 ms
 (slice1)   Executor memory: 51302K bytes avg x 96 workers, 52481K bytes max (seg17).
 (slice2)   Executor memory: 118204K bytes avg x 96 workers, 121890K bytes max (seg33).
Total runtime: 41892.667 ms
Оба Seq Scan сопровождаются Redistribute Motion — обе таблицы распределены не по customer_id. fact_orders создавалась с DISTRIBUTED BY order_id, dim_customer_scoring — с DISTRIBUTED BY scoring_id: оба раза выбирался суррогатный ключ по умолчанию, без оглядки на будущие JOIN-ы.

Исправление ключа дистрибуции (DISTRIBUTED BY)

-- смотрим текущий ключ дистрибуции обеих таблиц
SELECT localoid::regclass AS table_name, attrnums
FROM gp_distribution_policy
WHERE localoid IN ('fact_orders'::regclass, 'dim_customer_scoring'::regclass);
 
-- меняем ключ дистрибуции таблицы скоринга на customer_id
ALTER TABLE dim_customer_scoring
  SET DISTRIBUTED BY (customer_id);
 
-- обновляем статистику после физического перераспределения данных
VACUUM ANALYZE dim_customer_scoring;
Таблицу фактов трогать не пришлось: fact_orders — центральная таблица DWH, которую джойнят по customer_id в десятках других отчётов. Смена DISTRIBUTED BY у большой таблицы — операция дорогая, поэтому сначала стоит проверить, нельзя ли обойтись правкой меньшей таблицы.

Результат: локальный JOIN

Gather Motion 96:1  (slice1; segments: 96)  (cost=0.00..612100.20 rows=1 width=48)
  (actual time=3241.552..3241.601 rows=1 loops=1)
  ->  HashAggregate  (cost=598200.10..598200.11 rows=1 width=48)
        (actual time=3198.220..3198.221 rows=1 loops=1)
        Group Key: fact.customer_id
        ->  Hash Join  (cost=210400.00..585310.00 rows=1310000 width=32)
              (actual time=612.204..2981.552 rows=1298742 loops=1)
              Hash Cond: (fact.customer_id = scoring.customer_id)
              ->  Seq Scan on fact_orders fact
                    (cost=0.00..250340.00 rows=8400000 width=24)
                    (actual time=0.010..1102.334 rows=8391204 loops=1)
              ->  Hash  (cost=92100.00..92100.00 rows=1310000 width=16)
                    (actual time=608.552..608.552 rows=1298742 loops=1)
                    ->  Seq Scan on dim_customer_scoring scoring
                          (cost=0.00..92100.00 rows=1310000 width=16)
                          (actual time=0.008..401.221 rows=1298742 loops=1)
Planner: Pivotal Optimizer (GPORCA)
Planning time: 18.204 ms
 (slice1)   Executor memory: 42108K bytes avg x 96 workers, 43012K bytes max (seg0).
Total runtime: 3298.114 ms
Между Seq Scan и Hash Join больше нет ни одного Motion: обе таблицы распределены по customer_id, строки с одинаковым ключом лежат на одном сегменте, JOIN выполняется локально, без единого байта по интерконнекту. Единственный оставшийся Motion — Gather наверху плана. Итог: 41,9 секунды против 3,3 — ускорение больше чем в 12 раз без единого индекса.

Влияние Data Skew (перекоса данных) на план запроса

Даже правильный ключ дистрибуции не спасает, если значения в нём распределены неравномерно — это data skew. Классический пример: DISTRIBUTED BY country_code, когда 80% строк приходится на одну-две страны. JOIN формально становится локальным, Motion-узлов нет, но один сегмент обрабатывает в разы больше строк, чем остальные, и именно он определяет время запроса.
Перекос не всегда виден из EXPLAIN ANALYZE напрямую — там показано суммарное число строк по всем сегментам. Проверить распределение по факту помогает прямой запрос по gp_segment_id:
SELECT gp_segment_id, count(*)
FROM fact_orders
GROUP BY gp_segment_id
ORDER BY count(*) DESC
LIMIT 10;
Если count на одном сегменте в разы больше, чем на остальных — это skew, и лечится он сменой ключа дистрибуции на более равномерный либо DISTRIBUTED RANDOMLY. Оба варианта требуют актуальной статистики, поэтому VACUUM ANALYZE после смены DISTRIBUTED BY — обязательный шаг, а не опция.

Аудит DWH от экспертов DB Serv

Разбор одного JOIN-а в статье занимает пару минут — разбор всего DWH с сотнями таблиц, где вперемешку устаревшие ключи дистрибуции, забытая статистика и перекос данных, вручную занимает недели. Специалисты администрирования и аудита Greenplum и Arenadata DB от DB Serv проверяют ключи дистрибуции таблиц кластера, находят проблемные Motion-узлы в регулярных запросах и настраивают мониторинг, который сигнализирует о новых проблемах до того, как они превратятся в 40-секундный отчёт для бизнеса.

Краткие выводы

  • Motion-узлы (Broadcast, Redistribute, Gather) — главная причина долгих JOIN-ов в Greenplum, а не отсутствие индексов.
  • Redistribute Motion на обеих сторонах JOIN почти всегда означает, что DISTRIBUTED BY не совпадает с ключом соединения.
  • Broadcast Motion оправдан только для небольших справочных таблиц — для больших это признак устаревшей статистики.
  • Расхождение оценки планировщика и факта по rows — сигнал к VACUUM ANALYZE.
  • Локальный JOIN без Motion возможен, только когда обе таблицы распределены по одному ключу соединения.
  • Data Skew не виден в EXPLAIN ANALYZE напрямую — его проверяют через gp_segment_id.

Частые вопросы по теме

Оптимизация запросов и аудит DWH
от экспертов DB Serv
Забытая статистика, устаревшие ключи дистрибуции и перекос данных могут замедлить выполнение запросов на десятки минут. Мы проанализируем узкие места и оптимизируем архитектуру так, чтобы ваши JOIN-ы выполнялись локально, без лишней нагрузки на интерконнект.
Оставить заявку

Эксперт ДБ-сервис

Наши топ-3 компетенции по Arenadata DB
Каждое из наших направлений создано для того, чтобы ваше хранилище данных работало на максимальной скорости, а бизнес развивался без сбоев и непредсказуемых рисков.
  • Глубокий анализ производительности вашего хранилища данных. Выявляем узкие места в архитектуре, проверяем «тяжелые» запросы и ключи дистрибуции. Предоставляем четкие рекомендации и пошаговый план оптимизации кластера.
    Подробнее
  • Комплексное техническое сопровождение кластера. Грамотно распределяем ресурсы (Resource Groups), обеспечиваем бесперебойную работу ночных ETL-загрузок и дневной аналитики.
    Подробнее
  • Применяем системный и прозрачный подход. Понятный процесс работы: от детального сбора метрик конфигурации до профилирования нагрузки. Внедряем лучшие инженерные практики для выхода на новый уровень надежности.
    Подробнее
Еще статьи по теме