Настройка Cost-Based Optimizer в StarRocks

Настройка Cost-Based Optimizer в StarRocks

Один и тот же 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. Он выполняет запрос и печатает фактическое число строк и время по каждому узлу.

Результирующий EXPLAIN analyze показывающий реальный цифры для StarRocks DB JOIN

Сверяем оценки CBO с фактом по ключевым узлам.

Explain analyze для работы с JOIN операциями в StarRocks db

Узел плана Оценка 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). Разобрать эти темы вживую можно на митапе Школы Больших Данных.

Референсные ссылки