Содержание
- Что такое Cost-Based Optimizer в StarRocks
- Почему статистика решает больше, чем движок
- Практика: сбор и сброс статистики
- Команды ANALYZE для ручного и автоматического сбора
- Как сбросить статистику для эксперимента
- Читаем план: EXPLAIN и EXPLAIN COSTS
- План до сбора статистики
- План после сбора статистики
- Что изменилось: сравнительная таблица
- Проверяем план на реальном выполнении через EXPLAIN ANALYZE
- Управление планировщиком через хинты
- Заключение
- Референсные ссылки
Один и тот же SQL-запрос в StarRocks может выполниться за секунды или упасть по памяти. Разница не в железе, а в плане, который построил оптимизатор. За выбор плана отвечает Cost-Based Optimizer, а его точность целиком зависит от статистики данных. В этой статье разберем, как StarRocks собирает статистику, как читать план через EXPLAIN и EXPLAIN COSTS, и покажем на реальном примере, как сбор статистики меняет порядок соединений и втрое снижает оценку стоимости запроса.
Материал продолжает наш разбор StarRocks и опирается на тот же стенд с базой shop из статьи про сравнение StarRocks и ClickHouse. Умение читать планы выполнения это ключевой навык дата-инженера уровня Middle и выше, поэтому все команды приводим с пояснением, что именно смотреть.
Что такое Cost-Based Optimizer в StarRocks
Планировщик запроса это компонент, который превращает текст SQL в дерево физических операций. У него почти всегда есть несколько вариантов: в каком порядке соединять таблицы, какую сторону класть в хэш-таблицу, рассылать данные широковещательно или перераспределять по сети. Каждый вариант стоит по-разному в терминах процессора, памяти и сети.
Cost-Based Optimizer, сокращенно CBO, выбирает вариант с наименьшей оценочной стоимостью. Чтобы оценка была осмысленной, оптимизатору нужно знать характеристики данных: сколько строк в таблице, сколько уникальных значений в колонке, как распределены значения, какова средняя ширина строки. Это и есть статистика. Без нее CBO подставляет грубые значения по умолчанию и легко ошибается в разы.
Оптимизатор StarRocks построен по модели Cascades и глубоко заточен под векторизованный движок. Он умеет переупорядочивать соединения, переписывать подзапросы, выносить общие табличные выражения и выбирать стратегию распределенного JOIN. Все эти решения он принимает на основе статистики, поэтому сбор статистики это не формальность, а условие того, что план будет хорошим.
Проектирование Online-хранилищ данных на StarRocks.
Код курса
STAR
Ближайшая дата курса
7 сентября, 2026
Продолжительность
24 ак.часов
Стоимость обучения
76 800
Почему статистика решает больше, чем движок
Здесь полезно сравнить подход StarRocks с ClickHouse. ClickHouse исторически опирается на правила: порядок соединения во многом задается тем, как написан запрос, а правая таблица уходит в хэш-таблицу целиком. Для сложных многотабличных JOIN это приводит к ручной работе, когда инженер сам подбирает порядок или выносит справочники в словари оперативной памяти через функции вроде dictGet, чтобы вообще избежать соединения.
StarRocks перекладывает эту работу на CBO. Инженер пишет запрос в естественном виде, а оптимизатор сам решает, как его выполнить дешевле. Но есть условие: у оптимизатора должна быть свежая статистика. Если статистики нет, StarRocks оценивает селективность фильтра по умолчанию, обычно как половину таблицы, и может выбрать неоптимальный план. Дальше мы увидим это на цифрах.
Практика: сбор и сброс статистики
Команды ANALYZE для ручного и автоматического сбора
StarRocks собирает два вида статистики. Базовая статистика это число строк, число уникальных значений NDV, доля null, минимум и максимум, средняя ширина строки. Гистограмма это распределение значений колонки с наиболее частыми значениями MCV, она нужна для точной оценки селективности фильтров по перекошенным данным.
По умолчанию StarRocks собирает базовую статистику автоматически. Ручной сбор запускается командой ANALYZE TABLE и полезен сразу после массовой загрузки данных, чтобы не ждать автосборщик.
-- протестировано для StarRocks 3.5.0 -- полный синхронный сбор базовой статистики ANALYZE TABLE orders; ANALYZE TABLE order_items; ANALYZE TABLE customers; ANALYZE TABLE products; -- гистограмма по колонкам фильтра, здесь селективность важнее всего ANALYZE TABLE orders UPDATE HISTOGRAM ON order_date, status;
Проверить, что и когда собрано, помогают системные представления.
SHOW STATS META; -- базовая статистика по таблицам SHOW HISTOGRAM META; -- собранные гистограммы SHOW ANALYZE STATUS; -- статус задач сбора
Как сбросить статистику для эксперимента
Чтобы увидеть план в состоянии без статистики, ее нужно удалить и запретить автосборщику собрать заново. Без отключения автосбора статистика вернется через минуту, и чистого состояния не получить.
-- отключаем автосбор на время эксперимента
ADMIN SET FRONTEND CONFIG ("enable_statistic_collect" = "false");
-- удаляем собранную статистику
USE shop;
DROP STATS orders;
DROP STATS order_items;
DROP STATS customers;
-- проверяем, что статистики нет
SHOW STATS META;
-- ... снимаем план "до", потом собираем статистику и снимаем план "после" ...
-- в конце возвращаем автосбор
ADMIN SET FRONTEND CONFIG ("enable_statistic_collect" = "true");
Полный сценарий с обоими прогонами лежит в файле cbo_demo.sql в репозитории курса. Дальше разбираем, что именно меняется в плане.
Читаем план: EXPLAIN и EXPLAIN COSTS
Для демонстрации возьмем компактный запрос на три таблицы, чтобы план помещался на экран. Он считает число позаказных строк по регионам за последний квартал среди оплаченных заказов.
SELECT c.region_id, count(*) AS cnt FROM order_items oi JOIN orders o ON oi.order_id = o.order_id JOIN customers c ON o.customer_id = c.customer_id WHERE o.order_date >= '2025-04-01' AND o.status = 'paid' GROUP BY c.region_id;
Команда EXPLAIN показывает дерево операций и порядок соединений. Команда EXPLAIN COSTS добавляет к каждому узлу оценку числа строк и статистику колонок, а вверху печатает суммарную оценку PLAN COST по процессору и памяти. Именно COSTS показывает, что оптимизатор думает о данных.
План до сбора статистики
Снимаем план сразу после сброса статистики. Компактный EXPLAIN до ANALYZE выглядит так.
Explain String
-------------------------------------------------------------------------------+
PLAN FRAGMENT 0
OUTPUT EXPRS:11: region_id 13: count
PARTITION: UNPARTITIONED
RESULT SINK
14:EXCHANGE
limit: 200
PLAN FRAGMENT 1
OUTPUT EXPRS:
PARTITION: HASH_PARTITIONED: 11: region_id
STREAM DATA SINK
EXCHANGE ID: 14
UNPARTITIONED
13:AGGREGATE (merge finalize)
output: count(13: count)
group by: 11: region_id
limit: 200
12:EXCHANGE
PLAN FRAGMENT 2
OUTPUT EXPRS:
colocate exec groups: ExecGroup{groupId=2, nodeIds=[0, 9, 10, 11]}
PARTITION: RANDOM
STREAM DATA SINK
EXCHANGE ID: 12
HASH_PARTITIONED: 11: region_id
11:AGGREGATE (update serialize)
STREAMING
output: count(*)
group by: 11: region_id
10:Project
<slot 11> : 11: region_id
9:HASH JOIN
join op: INNER JOIN (BUCKET_SHUFFLE)
colocate: false, reason:
equal join conjunct: 2: order_id = 5: order_id
----8:EXCHANGE
0:OlapScanNode
TABLE: order_items
PREAGGREGATION: ON
PREDICATES: 2: order_id IS NOT NULL
partitions=1/1
rollup: order_items
tabletRatio=96/96
tabletList=13229,13231,13233,13235,13237,13239,13241,13243,13245,13247 ...
cardinality=360000000
avgRowSize=1.0
PLAN FRAGMENT 3
OUTPUT EXPRS:
PARTITION: HASH_PARTITIONED: 6: customer_id
STREAM DATA SINK
EXCHANGE ID: 08
BUCKET_SHUFFLE_HASH_PARTITIONED: 5: order_id
7:Project
<slot 5> : 5: order_id
<slot 11> : 11: region_id
6:HASH JOIN
join op: INNER JOIN (PARTITIONED)
colocate: false, reason:
equal join conjunct: 6: customer_id = 9: customer_id
----5:EXCHANGE
3:EXCHANGE
PLAN FRAGMENT 4
OUTPUT EXPRS:
PARTITION: RANDOM
STREAM DATA SINK
EXCHANGE ID: 05
HASH_PARTITIONED: 9: customer_id
4:OlapScanNode
TABLE: customers
PREAGGREGATION: ON
PREDICATES: 9: customer_id IS NOT NULL
partitions=1/1
rollup: customers
tabletRatio=16/16
tabletList=13161,13163,13165,13167,13169,13171,13173,13175,13177,13179 ...
cardinality=18000000
avgRowSize=2.0
PLAN FRAGMENT 5
OUTPUT EXPRS:
PARTITION: RANDOM
STREAM DATA SINK
EXCHANGE ID: 03
HASH_PARTITIONED: 6: customer_id
2:Project
<slot 5> : 5: order_id
<slot 6> : 6: customer_id
1:OlapScanNode
TABLE: orders
PREAGGREGATION: ON
PREDICATES: 7: order_date >= '2025-04-01', 8: status = 'paid'
partitions=1/1
rollup: orders
tabletRatio=48/48
tabletList=13063,13065,13067,13069,13071,13073,13075,13077,13079,13081 ...
cardinality=80000000
avgRowSize=4.0
Обратите внимание на узлы сканирования. Таблица orders после фильтра по дате и статусу оценена в 80 млн строк, ровно половину от 160 млн. Это значение по умолчанию, потому что без гистограммы StarRocks не знает реальную селективность и берет коэффициент 0.5. Первое соединение идет как orders с customers по стратегии PARTITIONED, а в хэш-таблицу уходит customers.
Полная версия с оценками стоимости через EXPLAIN COSTS до ANALYZE показательнее.
Explain String
------------------------------------------------------------------------------------------------------------------------+
PLAN COST
CPU: 2.63040024E10
Memory: 1.3920024E9
PLAN FRAGMENT 0(F08)
Output Exprs:11: region_id 13: count
Input Partition: UNPARTITIONED
RESULT SINK
14:EXCHANGE
distribution type: GATHER
limit: 200
cardinality: 200
PLAN FRAGMENT 1(F07)
Input Partition: HASH_PARTITIONED: 11: region_id
OutPut Partition: UNPARTITIONED
OutPut Exchange Id: 14
13:AGGREGATE (merge finalize)
aggregate: count[([13: count, BIGINT, false]); args: ; result: BIGINT; args nullable: true; result nullable: false]
group by: [11: region_id, INT, true]
limit: 200
cardinality: 200
column statistics:
* region_id-->[-Infinity, Infinity, 0.0, 1.0, 1.0] UNKNOWN
* count-->[-Infinity, Infinity, 0.0, 8.0, 1.0] ESTIMATE
12:EXCHANGE
distribution type: SHUFFLE
partition exprs: [11: region_id, INT, true]
cardinality: 180000000
PLAN FRAGMENT 2(F00)
Input Partition: RANDOM
OutPut Partition: HASH_PARTITIONED: 11: region_id
OutPut Exchange Id: 12
11:AGGREGATE (update serialize)
STREAMING
aggregate: count[(*); args: ; result: BIGINT; args nullable: false; result nullable: false]
group by: [11: region_id, INT, true]
cardinality: 180000000
column statistics:
* region_id-->[-Infinity, Infinity, 0.0, 1.0, 1.0] UNKNOWN
* count-->[-Infinity, Infinity, 0.0, 8.0, 1.0] ESTIMATE
10:Project
output columns:
11 <-> [11: region_id, INT, true]
cardinality: 360000000
column statistics:
* region_id-->[-Infinity, Infinity, 0.0, 1.0, 1.0] UNKNOWN
9:HASH JOIN
join op: INNER JOIN (BUCKET_SHUFFLE)
equal join conjunct: [2: order_id, BIGINT, true] = [5: order_id, BIGINT, true]
build runtime filters:
- filter_id = 1, build_expr = (5: order_id), remote = false
output columns: 11
can local shuffle: true
cardinality: 360000000
column statistics:
* customer_id-->[-Infinity, Infinity, 0.0, 1.0, 1.0] UNKNOWN
* customer_id-->[-Infinity, Infinity, 0.0, 1.0, 1.0] UNKNOWN
* region_id-->[-Infinity, Infinity, 0.0, 1.0, 1.0] UNKNOWN
----8:EXCHANGE
distribution type: SHUFFLE
partition exprs: [5: order_id, BIGINT, true]
cardinality: 80000000
0:OlapScanNode
table: order_items, rollup: order_items
preAggregation: on
Predicates: 2: order_id IS NOT NULL
partitionsRatio=1/1, tabletsRatio=96/96
tabletList=13229,13231,13233,13235,13237,13239,13241,13243,13245,13247 ...
actualRows=400000000, avgRowSize=1.0
cardinality: 360000000
probe runtime filters:
- filter_id = 1, probe_expr = (2: order_id)
column statistics:
* order_id-->[-Infinity, Infinity, 0.0, 1.0, 1.0] UNKNOWN
PLAN FRAGMENT 3(F05)
Input Partition: HASH_PARTITIONED: 6: customer_id
OutPut Partition: BUCKET_SHUFFLE_HASH_PARTITIONED: 5: order_id
OutPut Exchange Id: 08
7:Project
output columns:
5 <-> [5: order_id, BIGINT, true]
11 <-> [11: region_id, INT, true]
cardinality: 80000000
column statistics:
* order_id-->[-Infinity, Infinity, 0.0, 1.0, 1.0] UNKNOWN
* region_id-->[-Infinity, Infinity, 0.0, 1.0, 1.0] UNKNOWN
6:HASH JOIN
join op: INNER JOIN (PARTITIONED)
equal join conjunct: [6: customer_id, BIGINT, true] = [9: customer_id, BIGINT, true]
build runtime filters:
- filter_id = 0, build_expr = (9: customer_id), remote = true
output columns: 5, 11
can local shuffle: true
cardinality: 80000000
column statistics:
* order_id-->[-Infinity, Infinity, 0.0, 1.0, 1.0] UNKNOWN
* customer_id-->[-Infinity, Infinity, 0.0, 1.0, 1.0] UNKNOWN
* customer_id-->[-Infinity, Infinity, 0.0, 1.0, 1.0] UNKNOWN
* region_id-->[-Infinity, Infinity, 0.0, 1.0, 1.0] UNKNOWN
----5:EXCHANGE
distribution type: SHUFFLE
partition exprs: [9: customer_id, BIGINT, true]
cardinality: 18000000
3:EXCHANGE
distribution type: SHUFFLE
partition exprs: [6: customer_id, BIGINT, true]
cardinality: 80000000
PLAN FRAGMENT 4(F03)
Input Partition: RANDOM
OutPut Partition: HASH_PARTITIONED: 9: customer_id
OutPut Exchange Id: 05
4:OlapScanNode
table: customers, rollup: customers
preAggregation: on
Predicates: 9: customer_id IS NOT NULL
partitionsRatio=1/1, tabletsRatio=16/16
tabletList=13161,13163,13165,13167,13169,13171,13173,13175,13177,13179 ...
actualRows=20000000, avgRowSize=2.0
cardinality: 18000000
column statistics:
* customer_id-->[-Infinity, Infinity, 0.0, 1.0, 1.0] UNKNOWN
* region_id-->[-Infinity, Infinity, 0.0, 1.0, 1.0] UNKNOWN
PLAN FRAGMENT 5(F01)
Input Partition: RANDOM
OutPut Partition: HASH_PARTITIONED: 6: customer_id
OutPut Exchange Id: 03
2:Project
output columns:
5 <-> [5: order_id, BIGINT, true]
6 <-> [6: customer_id, BIGINT, true]
cardinality: 80000000
column statistics:
* order_id-->[-Infinity, Infinity, 0.0, 1.0, 1.0] UNKNOWN
* customer_id-->[-Infinity, Infinity, 0.0, 1.0, 1.0] UNKNOWN
1:OlapScanNode
table: orders, rollup: orders
preAggregation: on
Predicates: [7: order_date, DATE, true] >= '2025-04-01', [8: status, VARCHAR(16), true] = 'paid'
partitionsRatio=1/1, tabletsRatio=48/48
tabletList=13063,13065,13067,13069,13071,13073,13075,13077,13079,13081 ...
actualRows=160000000, avgRowSize=4.0
cardinality: 80000000
probe runtime filters:
- filter_id = 0, probe_expr = (6: customer_id)
column statistics:
* order_id-->[-Infinity, Infinity, 0.0, 1.0, 1.0] UNKNOWN
* customer_id-->[-Infinity, Infinity, 0.0, 1.0, 1.0] UNKNOWN
* order_date-->[1.7434656E9, Infinity, 0.0, 1.0, 1.0] UNKNOWN
* status-->[-Infinity, Infinity, 0.0, 1.0, 1.0] UNKNOWN
Здесь видно два важных признака отсутствия статистики. Во-первых, у всех колонок в блоках column statistics стоит пометка UNKNOWN, оптимизатор ничего не знает о распределении значений. Во-вторых, итоговая оценка результата раздута до 180 млн строк, а суммарная стоимость PLAN COST по процессору равна 2.63 на 10 в степени 10. Средняя ширина строки orders взята как 4 байта, тоже заглушка.
План после сбора статистики
Теперь собираем базовую статистику и гистограммы, затем снимаем тот же план. Компактный EXPLAIN после ANALYZE меняется.
Explain String
-------------------------------------------------------------------------------+
PLAN FRAGMENT 0
OUTPUT EXPRS:11: region_id 13: count
PARTITION: UNPARTITIONED
RESULT SINK
13:EXCHANGE
limit: 200
PLAN FRAGMENT 1
OUTPUT EXPRS:
PARTITION: HASH_PARTITIONED: 11: region_id
STREAM DATA SINK
EXCHANGE ID: 13
UNPARTITIONED
12:AGGREGATE (merge finalize)
output: count(13: count)
group by: 11: region_id
limit: 200
11:EXCHANGE
PLAN FRAGMENT 2
OUTPUT EXPRS:
colocate exec groups: ExecGroup{groupId=2, nodeIds=[0, 8, 9, 10]}
PARTITION: RANDOM
STREAM DATA SINK
EXCHANGE ID: 11
HASH_PARTITIONED: 11: region_id
10:AGGREGATE (update serialize)
STREAMING
output: count(*)
group by: 11: region_id
9:Project
<slot 11> : 11: region_id
8:HASH JOIN
join op: INNER JOIN (BUCKET_SHUFFLE)
colocate: false, reason:
equal join conjunct: 2: order_id = 5: order_id
----7:EXCHANGE
0:OlapScanNode
TABLE: order_items
PREAGGREGATION: ON
PREDICATES: 2: order_id IS NOT NULL
partitions=1/1
rollup: order_items
tabletRatio=96/96
tabletList=13229,13231,13233,13235,13237,13239,13241,13243,13245,13247 ...|
cardinality=400000000
avgRowSize=8.0
PLAN FRAGMENT 3
OUTPUT EXPRS:
PARTITION: RANDOM
STREAM DATA SINK
EXCHANGE ID: 07
BUCKET_SHUFFLE_HASH_PARTITIONED: 5: order_id
6:Project
<slot 5> : 5: order_id
<slot 11> : 11: region_id
5:HASH JOIN
join op: INNER JOIN (BUCKET_SHUFFLE)
colocate: false, reason:
equal join conjunct: 9: customer_id = 6: customer_id
----4:EXCHANGE
1:OlapScanNode
TABLE: customers
PREAGGREGATION: ON
PREDICATES: 9: customer_id IS NOT NULL
partitions=1/1
rollup: customers
tabletRatio=16/16
tabletList=13161,13163,13165,13167,13169,13171,13173,13175,13177,13179 ...|
cardinality=20000000
avgRowSize=12.0
PLAN FRAGMENT 4
OUTPUT EXPRS:
PARTITION: RANDOM
STREAM DATA SINK
EXCHANGE ID: 04
BUCKET_SHUFFLE_HASH_PARTITIONED: 6: customer_id
3:Project
<slot 5> : 5: order_id
<slot 6> : 6: customer_id
2:OlapScanNode
TABLE: orders
PREAGGREGATION: ON
PREDICATES: 7: order_date >= '2025-04-01', 8: status = 'paid'
partitions=1/1
rollup: orders
tabletRatio=48/48
tabletList=13063,13065,13067,13069,13071,13073,13075,13077,13079,13081 ...|
cardinality=9154400
avgRowSize=26.666664
Порядок соединений переставился. Теперь orders после фильтра оценен в 9.15 млн строк, и именно он, как самая маленькая сторона, уходит в build первого соединения. Первое соединение сменило стратегию с PARTITIONED на BUCKET_SHUFFLE, которое дешевле, потому что использует бакетирование таблиц. Огромный order_items по-прежнему течет пробной стороной, но хэш-таблица, через которую он проходит, стала в разы меньше.
Версия EXPLAIN COSTS после ANALYZE подтверждает разницу цифрами.
Explain String
-------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
PLAN COST
CPU: 8.440992867476382E9
Memory: 2.5632438525424263E8
PLAN FRAGMENT 0(F06)
Output Exprs:11: region_id 13: count
Input Partition: UNPARTITIONED
RESULT SINK
13:EXCHANGE
distribution type: GATHER
limit: 200
cardinality: 90
PLAN FRAGMENT 1(F05)
Input Partition: HASH_PARTITIONED: 11: region_id
OutPut Partition: UNPARTITIONED
OutPut Exchange Id: 13
12:AGGREGATE (merge finalize)
aggregate: count[([13: count, BIGINT, false]); args: ; result: BIGINT; args nullable: true; result nullable: false]
group by: [11: region_id, INT, true]
limit: 200
cardinality: 90
column statistics:
* region_id-->[1.0, 90.0, 0.0, 4.0, 90.0] ESTIMATE
* count-->[0.0, 2.490471357548823E7, 0.0, 8.0, 90.0] ESTIMATE
11:EXCHANGE
distribution type: SHUFFLE
partition exprs: [11: region_id, INT, true]
cardinality: 90
PLAN FRAGMENT 2(F00)
Input Partition: RANDOM
OutPut Partition: HASH_PARTITIONED: 11: region_id
OutPut Exchange Id: 11
10:AGGREGATE (update serialize)
STREAMING
aggregate: count[(*); args: ; result: BIGINT; args nullable: false; result nullable: false]
group by: [11: region_id, INT, true]
cardinality: 90
column statistics:
* region_id-->[1.0, 90.0, 0.0, 4.0, 90.0] ESTIMATE
* count-->[0.0, 2.490471357548823E7, 0.0, 8.0, 90.0] ESTIMATE
9:Project
output columns:
11 <-> [11: region_id, INT, true]
cardinality: 24904714
column statistics:
* region_id-->[1.0, 90.0, 0.0, 4.0, 90.0] ESTIMATE
8:HASH JOIN
join op: INNER JOIN (BUCKET_SHUFFLE)
equal join conjunct: [2: order_id, BIGINT, true] = [5: order_id, BIGINT, true]
build runtime filters:
- filter_id = 1, build_expr = (5: order_id), remote = false
output columns: 11
can local shuffle: true
cardinality: 24904714
column statistics:
* order_id-->[1.0, 1.59999999E8, 0.0, 8.0, 9154399.901937237] ESTIMATE
* order_id-->[1.0, 1.59999999E8, 0.0, 8.0, 9154399.901937237] ESTIMATE
* region_id-->[1.0, 90.0, 0.0, 4.0, 90.0] ESTIMATE
----7:EXCHANGE
distribution type: SHUFFLE
partition exprs: [5: order_id, BIGINT, true]
cardinality: 9154400
0:OlapScanNode
table: order_items, rollup: order_items
preAggregation: on
Predicates: 2: order_id IS NOT NULL
partitionsRatio=1/1, tabletsRatio=96/96
tabletList=13229,13231,13233,13235,13237,13239,13241,13243,13245,13247 ...
actualRows=400000000, avgRowSize=8.0
cardinality: 400000000
probe runtime filters:
- filter_id = 1, probe_expr = (2: order_id)
column statistics:
* order_id-->[1.0, 1.59999999E8, 0.0, 8.0, 1.470308E8] ESTIMATE
PLAN FRAGMENT 3(F01)
Input Partition: RANDOM
OutPut Partition: BUCKET_SHUFFLE_HASH_PARTITIONED: 5: order_id
OutPut Exchange Id: 07
6:Project
output columns:
5 <-> [5: order_id, BIGINT, true]
11 <-> [11: region_id, INT, true]
cardinality: 9154400
column statistics:
* order_id-->[1.0, 1.6E8, 0.0, 8.0, 9154399.901937237] ESTIMATE
* region_id-->[1.0, 90.0, 0.0, 4.0, 90.0] ESTIMATE
5:HASH JOIN
join op: INNER JOIN (BUCKET_SHUFFLE)
equal join conjunct: [9: customer_id, BIGINT, true] = [6: customer_id, BIGINT, true]
build runtime filters:
- filter_id = 0, build_expr = (6: customer_id), remote = false
output columns: 5, 11
can local shuffle: true
cardinality: 9154400
column statistics:
* order_id-->[1.0, 1.6E8, 0.0, 8.0, 9154399.901937237] ESTIMATE
* customer_id-->[1.0, 2.0E7, 0.0, 8.0, 9154399.901937237] ESTIMATE
* customer_id-->[1.0, 2.0E7, 0.0, 8.0, 9154399.901937237] ESTIMATE
* region_id-->[1.0, 90.0, 0.0, 4.0, 90.0] ESTIMATE
----4:EXCHANGE
distribution type: SHUFFLE
partition exprs: [6: customer_id, BIGINT, true]
cardinality: 9154400
1:OlapScanNode
table: customers, rollup: customers
preAggregation: on
Predicates: 9: customer_id IS NOT NULL
partitionsRatio=1/1, tabletsRatio=16/16
tabletList=13161,13163,13165,13167,13169,13171,13173,13175,13177,13179 ...
actualRows=20000000, avgRowSize=12.0
cardinality: 20000000
probe runtime filters:
- filter_id = 0, probe_expr = (9: customer_id)
column statistics:
* customer_id-->[1.0, 2.0E7, 0.0, 8.0, 2.0E7] ESTIMATE
* region_id-->[1.0, 90.0, 0.0, 4.0, 90.0] ESTIMATE
PLAN FRAGMENT 4(F02)
Input Partition: RANDOM
OutPut Partition: BUCKET_SHUFFLE_HASH_PARTITIONED: 6: customer_id
OutPut Exchange Id: 04
3:Project
output columns:
5 <-> [5: order_id, BIGINT, true]
6 <-> [6: customer_id, BIGINT, true]
cardinality: 9154400
column statistics:
* order_id-->[1.0, 1.6E8, 0.0, 8.0, 9154399.901937237] ESTIMATE
* customer_id-->[1.0, 2.0E7, 0.0, 8.0, 9154399.901937237] ESTIMATE
2:OlapScanNode
table: orders, rollup: orders
preAggregation: on
Predicates: [7: order_date, DATE, true] >= '2025-04-01', DictDecode(14: status, [<place-holder> = 'paid'])
dict_col=status
partitionsRatio=1/1, tabletsRatio=48/48
tabletList=13063,13065,13067,13069,13071,13073,13075,13077,13079,13081 ...
actualRows=160000000, avgRowSize=26.666664
cardinality: 9154400
column statistics:
* order_id-->[1.0, 1.6E8, 0.0, 8.0, 9154399.901937237] ESTIMATE
* customer_id-->[1.0, 2.0E7, 0.0, 8.0, 9154399.901937237] ESTIMATE
* order_date-->[1.7434656E9, 1.7515872E9, 0.0, 4.0, 553.0] MCV: [[2025-05-25:259344][2025-06-14:258768][2025-06-29:258656][2025-05-14:257408][2025-06-08:257136]] ESTIMATE|
* status-->[-Infinity, Infinity, 0.0, 6.6666643, 3.0] MCV: [[paid:52410224]] ESTIMATE
Колонки теперь помечены ESTIMATE, у каждой есть NDV и границы значений. По колонке status собрана гистограмма с наиболее частым значением paid, а по order_date виден список частых дат. StarRocks дополнительно применил словарное кодирование к status, это видно по DictDecode в предикате, что ускоряет фильтрацию низкокардинальной колонки. Итоговая оценка результата упала до 90 строк, ровно по числу регионов, а суммарная стоимость PLAN COST по процессору снизилась до 8.44 на 10 в степени 9.
Проектирование Online-хранилищ данных на StarRocks.
Код курса
STAR
Ближайшая дата курса
7 сентября, 2026
Продолжительность
24 ак.часов
Стоимость обучения
76 800
Что изменилось: сравнительная таблица
Соберем ключевые метрики обоих планов в одну таблицу. Все числа взяты из реальных выводов EXPLAIN COSTS на стенде.
| Показатель | До ANALYZE | После ANALYZE |
|---|---|---|
| Оценка строк orders после фильтра | 80 000 000 | 9 154 400 |
| Средняя ширина строки orders | 4.0, заглушка | 26.67, реальная |
| Статистика колонок | UNKNOWN | ESTIMATE с NDV и MCV |
| Первое соединение | orders с customers, build customers 18 млн, PARTITIONED | customers с orders, build orders 9.15 млн, BUCKET_SHUFFLE |
| Кардинальность верхнего JOIN с order_items | 360 000 000 | 24 904 714 |
| Итоговая оценка результата | 180 000 000 | 90 |
| PLAN COST, процессор | 2.63 на 10 в степени 10 | 8.44 на 10 в степени 9 |
| PLAN COST, память | 1.39 ГБ | 256 МБ |
Вывод простой. Единственное, что мы сделали, это собрали статистику. Ни строчки в запросе не изменилось. А оптимизатор пересчитал селективность фильтра с половины таблицы до реальных 5.7%, переставил порядок соединений, сменил стратегию распределения и втрое снизил оценку стоимости. Это и есть работа CBO, ради которой статистику держат свежей.
Проверяем план на реальном выполнении через EXPLAIN ANALYZE
EXPLAIN COSTS показывает только оценки. Чтобы убедиться, что после сбора статистики оптимизатор не просто пересчитал числа, но и попал в реальность, запускаем EXPLAIN ANALYZE. Он выполняет запрос и печатает фактическое число строк и время по каждому узлу.
Сверяем оценки CBO с фактом по ключевым узлам.
| Узел плана | Оценка CBO | Факт выполнения |
|---|---|---|
| Скан orders после фильтра | 9 154 399 | 9 192 460 |
| Верхнее соединение с order_items | 24 904 713 | 22 977 811 |
| Итоговый результат по регионам | 90 | 90 |
Оценки совпали с реальностью почти точно, это и есть признак свежей статистики. Общее время около 6.3 секунды, пиковая память на инстанс держалась около 370 МБ. Самый тяжелый узел здесь не соединение, а скан orders, почти 30% времени и 914 миллисекунд. Причина в том, что фильтр идет по order_date и status, а ключ сортировки таблицы это order_id. Движку приходится читать все 160 млн строк, чтобы отобрать 9.19 млн. Это прямая подсказка к следующему шагу оптимизации: партиционировать orders по дате, тогда скан станет избирательным. Но это уже тема отдельной статьи про модель хранения.
Управление планировщиком через хинты
В большинстве случаев CBO выбирает план сам, и вмешиваться не нужно. Но иногда данные меняются быстрее, чем обновляется статистика, или инженер точно знает лучшую стратегию. Тогда помогают хинты, которые принудительно задают поведение планировщика.
Хинт стратегии соединения ставится в квадратных скобках сразу после ключевого слова JOIN. Он заставляет StarRocks использовать конкретный способ распределения независимо от оценок.
-- принудительный broadcast маленькой таблицы на все узлы SELECT c.region_id, count(*) FROM orders o JOIN [broadcast] customers c ON o.customer_id = c.customer_id GROUP BY c.region_id; -- принудительный shuffle для двух больших таблиц SELECT ... FROM order_items oi JOIN [shuffle] orders o ON oi.order_id = o.order_id;
Доступны стратегии broadcast, shuffle, bucket и colocate. Отдельно системные переменные задаются через хинт SET_VAR прямо в запросе, например чтобы поднять лимит памяти только для тяжелого запроса.
SELECT /*+ SET_VAR(query_mem_limit = 8589934592) */
c.region_id, count(*)
FROM order_items oi
JOIN orders o ON oi.order_id = o.order_id
JOIN customers c ON o.customer_id = c.customer_id
WHERE o.order_date >= '2025-04-01' AND o.status = 'paid'
GROUP BY c.region_id;
Хинты это точечный инструмент. Прежде чем ими пользоваться, всегда стоит убедиться, что дело не в устаревшей статистике. В девяти случаях из десяти свежий ANALYZE решает проблему без ручного вмешательства.
Заключение
Cost-Based Optimizer в StarRocks превращает статистику в скорость. На нашем примере сбор статистики без единой правки запроса переставил порядок соединений, сменил стратегию распределения и втрое снизил оценку стоимости плана. Практический вывод для дата-инженера такой: собирайте статистику после массовых загрузок, стройте гистограммы по колонкам фильтров и читайте план через EXPLAIN COSTS, а не гадайте. Умение находить в плане узкие места по оценкам строк отличает инженера уровня Middle и выше.
Мы разбираем десятки реальных планов выполнения на курсе Проектирование и разработка Online-хранилищ данных на StarRocks (код курса STAR). Навыки профилирования запросов и работы со статистикой полезны и для другой аналитической СУБД, поэтому логично закрепить их и на курсе по ClickHouse (код курса CLICH). Разобрать эти темы вживую можно на митапе Школы Больших Данных.




