Аналитические запросы: GROUPING SETS, ROLLUP, CUBE без UDF

GROUPING SETS, ROLLUP, CUBE - многомерная агрегация за один проход. GROUPING_ID для идентификации уровней, физический план Expand, сравнение с UNION ALL.

core

Проблема: BI-витрина с множеством разрезов

Представьте задачу: построить отчёт о продажах e-commerce платформы. Бизнесу нужны цифры в нескольких разрезах одновременно:

  • Итого по всей компании (Grand Total)
  • По каждой стране
  • По стране и категории товара
  • По стране, категории и каналу продаж

Наивный подход - написать четыре отдельных groupBy и объединить через UNION ALL:

from pyspark.sql import functions as F

df = spark.read.parquet("s3://bucket/sales/")

# Наивный подход: 4 отдельных группировки
grand_total = df.agg(F.sum("revenue").alias("revenue")) \
    .withColumn("country",  F.lit(None)) \
    .withColumn("category", F.lit(None)) \
    .withColumn("channel",  F.lit(None))

by_country = df.groupBy("country") \
    .agg(F.sum("revenue").alias("revenue")) \
    .withColumn("category", F.lit(None)) \
    .withColumn("channel",  F.lit(None))

by_country_cat = df.groupBy("country", "category") \
    .agg(F.sum("revenue").alias("revenue")) \
    .withColumn("channel", F.lit(None))

by_all = df.groupBy("country", "category", "channel") \
    .agg(F.sum("revenue").alias("revenue"))

result = grand_total \
    .unionByName(by_country) \
    .unionByName(by_country_cat) \
    .unionByName(by_all)

Этот код работает, но ужасно неэффективен. Каждый groupBy - это отдельный Spark Job со своим Shuffle. Данные читаются с диска четыре раза. При 50 GB исходных данных это 200 GB I/O вместо 50 GB. При добавлении нового разреза нужно добавить ещё один блок кода и ещё один полный scan.

Решение: операторы многомерной агрегации - GROUPING SETS, ROLLUP и CUBE. Они позволяют вычислить все нужные комбинации за один проход по данным с одним Shuffle.


Проблема антипаттерна UNION ALL

Прежде чем разбирать решение - зафиксируем точную стоимость антипаттерна. На диаграмме видно, почему UNION ALL масштабируется плохо:

При N разрезах UNION ALL требует N полных сканирований. GROUPING SETS - всегда один scan, независимо от числа разрезов. Разница становится принципиальной при N = 5–10.


GROUPING SETS: явный список комбинаций

GROUPING SETS - наиболее гибкий оператор. Вы явно перечисляете нужные комбинации группировок. Spark вычислит агрегаты только для указанных комбинаций, пропустив остальные.

Синтаксис в Spark SQL:

SELECT
    country,
    category,
    channel,
    SUM(revenue) AS total_revenue,
    COUNT(*)     AS order_count
FROM sales
GROUP BY GROUPING SETS (
    (),                          -- Grand Total: итого по всей компании
    (country),                   -- Итого по каждой стране
    (country, category),         -- Итого по стране + категории
    (country, category, channel) -- Детальный уровень
)

Эквивалент в DataFrame API:

from pyspark.sql import functions as F

df.groupBy(
    F.grouping_sets(
        [],
        ["country"],
        ["country", "category"],
        ["country", "category", "channel"]
    )
).agg(
    F.sum("revenue").alias("total_revenue"),
    F.count("*").alias("order_count")
)

Результат для трёх стран, двух категорий, двух каналов:

country | category    | channel  | total_revenue | order_count
NULL    | NULL        | NULL     | 1_250_000     | 48_500   ← Grand Total
RU      | NULL        | NULL     |   650_000     | 25_200   ← Итого Россия
DE      | NULL        | NULL     |   380_000     | 15_100   ← Итого Германия
FR      | NULL        | NULL     |   220_000     |  8_200   ← Итого Франция
RU      | Electronics | NULL     |   420_000     | 12_000   ← RU + Electronics
RU      | Clothing    | NULL     |   230_000     | 13_200   ← RU + Clothing
...
RU      | Electronics | Online   |   290_000     |  8_500   ← Детальный уровень
RU      | Electronics | Offline  |   130_000     |  3_500
...

NULL в колонке означает «агрегировано по всем значениям этой колонки». Именно здесь возникает первая проблема: как отличить технический NULL (агрегированная строка) от реального NULL (отсутствие данных в источнике)?


ROLLUP: иерархическая агрегация

ROLLUP - это специальный случай GROUPING SETS для иерархических данных. Он автоматически генерирует все «префиксные» комбинации измерений, двигаясь от самого детального к самому общему уровню.

ROLLUP(a, b, c) эквивалентен:

GROUP BY GROUPING SETS (
    (a, b, c),   -- полный детальный уровень
    (a, b),      -- промежуточный итог по a+b
    (a),         -- промежуточный итог по a
    ()           -- Grand Total
)

Число уровней = число измерений + 1 (для Grand Total). ROLLUP(year, month, day) даст 4 уровня: год/месяц/день, год/месяц, год, Grand Total.

Синтаксис:

SELECT
    year,
    month,
    day,
    SUM(revenue) AS revenue
FROM sales
GROUP BY ROLLUP (year, month, day)
ORDER BY year, month, day
# DataFrame API
df.rollup("year", "month", "day") \
  .agg(F.sum("revenue").alias("revenue")) \
  .orderBy("year", "month", "day")

Результат для 2024 года, Q1:

year | month | day  | revenue
2024 | 1     | 1    | 45_000   ← детальный день
2024 | 1     | 2    | 52_000
...
2024 | 1     | NULL | 850_000  ← итого за январь
2024 | 2     | NULL | 920_000  ← итого за февраль
2024 | 3     | NULL | 780_000  ← итого за март
2024 | NULL  | NULL | 2_550_000 ← итого за 2024 год
NULL | NULL  | NULL | 2_550_000 ← Grand Total

Ключевое свойство ROLLUP - иерархичность. Комбинация (month, NULL) встречается только в контексте каждого конкретного year. Нельзя получить итог по месяцу без привязки к году. Это отличает ROLLUP от CUBE.

Типичные use cases ROLLUP:

  • Финансовые отчёты: Компания → Дивизион → Отдел → Сотрудник
  • Временные ряды: Год → Квартал → Месяц → Неделя → День
  • Географические иерархии: Страна → Регион → Город → Район
-- Организационная иерархия
SELECT
    company, division, department, employee,
    SUM(salary) AS total_payroll
FROM employees
GROUP BY ROLLUP (company, division, department, employee)

CUBE: полная многомерная агрегация

CUBE генерирует все возможные комбинации измерений, включая все подмножества. Для N измерений это 2^N комбинаций.

CUBE(a, b, c) эквивалентен:

GROUP BY GROUPING SETS (
    (a, b, c),   -- 3 измерения
    (a, b),      -- 2 измерения
    (a, c),
    (b, c),
    (a),         -- 1 измерение
    (b),
    (c),
    ()           -- Grand Total
)

8 комбинаций для 3 измерений (2³ = 8). Для 4 измерений - 16, для 5 - 32. Это экспоненциальный рост.

Синтаксис:

SELECT
    country,
    category,
    channel,
    SUM(revenue) AS revenue,
    AVG(order_value) AS avg_check
FROM sales
GROUP BY CUBE (country, category, channel)
# DataFrame API
df.cube("country", "category", "channel") \
  .agg(
      F.sum("revenue").alias("revenue"),
      F.avg("order_value").alias("avg_check")
  )

Отличие от ROLLUP: CUBE включает «поперечные» комбинации. Например, (category, channel) - итого по категории вне зависимости от страны. В ROLLUP такой комбинации нет - только (country, category).

Когда CUBE необходим: когда все измерения равноправны и нет заданной иерархии. Например, OLAP-куб для анализа продаж по трём независимым осям: страна, категория, канал.

Предупреждение об экспоненциальном росте:

Измерений Комбинаций Строк результата (1000 уникальных значений/измерение)
2 4 ~2 004
3 8 ~3 008
4 16 ~16 016
5 32 ~5 032
6 64 ~64 064
10 1024 ~10 240 + overhead

На высококардинальных измерениях (тысячи уникальных значений) CUBE на 5–6 колонках может дать миллионы строк результата. Это нагружает как Shuffle, так и BI-инструмент.

Правило безопасного CUBE: не более 4–5 измерений с умеренной кардинальностью (< 100 уникальных значений). Для большего числа измерений используйте GROUPING SETS с явным списком нужных комбинаций.


Сравнение операторов

Оператор Комбинации Когда использовать
GROUPING SETS Явный список Точно знаете нужные разрезы, максимальная гибкость
ROLLUP N+1 (иерархия) Временные ряды, организационные иерархии
CUBE 2^N (все) OLAP-кубы, равноправные измерения, <= 4-5 колонок

GROUPING и GROUPING_ID: идентификация уровней

Вернёмся к проблеме NULL. После выполнения ROLLUP/CUBE в результате появляются строки с NULL в ряде колонок. Это создаёт неоднозначность: NULL может означать либо «агрегировано по всем значениям», либо «в исходных данных была запись с NULL».

-- В исходных данных есть строки с country = NULL (неизвестная страна)
-- После ROLLUP: как отличить итоговую строку от строки с "неизвестной страной"?
SELECT country, SUM(revenue)
FROM sales
GROUP BY ROLLUP(country)
-- Будет ДВЕ строки с country = NULL:
-- 1. Агрегированный Grand Total (технический NULL от ROLLUP)
-- 2. Итого по строкам с country = NULL в источнике (реальный NULL)

Для решения этой проблемы существуют функции GROUPING() и GROUPING_ID().

GROUPING(): индикатор для одной колонки

GROUPING(col) возвращает 1 если данная колонка агрегирована (NULL создан ROLLUP/CUBE), 0 если колонка реально участвует в группировке:

SELECT
    country,
    category,
    GROUPING(country)  AS country_is_aggregated,
    GROUPING(category) AS category_is_aggregated,
    SUM(revenue) AS revenue
FROM sales
GROUP BY CUBE (country, category)
country | category    | country_is_agg | cat_is_agg | revenue
RU      | Electronics | 0              | 0          | 420_000  ← детальный
RU      | NULL        | 0              | 1          | 650_000  ← итого по RU (все категории)
NULL    | Electronics | 1              | 0          | 680_000  ← итого по Electronics (все страны)
NULL    | NULL        | 1              | 1          | 1_250_000 ← Grand Total

Теперь можно точно различить строки. Строка с country=NULL, country_is_agg=1 - это Grand Total. Строка с country=NULL, country_is_agg=0 - это реальные данные с неизвестной страной.

GROUPING_ID(): битовая маска уровня

GROUPING_ID(col1, col2, ...) возвращает целое число - битовую маску, где каждый бит соответствует одной колонке. Бит = 1, если колонка агрегирована; бит = 0, если колонка реально участвует в группировке.

Порядок битов: правый бит = последняя колонка в списке:

SELECT
    country,
    category,
    channel,
    GROUPING_ID(country, category, channel) AS level_id,
    SUM(revenue) AS revenue
FROM sales
GROUP BY CUBE (country, category, channel)
country | category    | channel | level_id | revenue   | Описание
RU      | Electronics | Online  | 0  (000) | 290_000   | детальный уровень
RU      | Electronics | NULL    | 1  (001) | 420_000   | country + category
RU      | NULL        | Online  | 2  (010) | 380_000   | country + channel
RU      | NULL        | NULL    | 3  (011) | 650_000   | только country
NULL    | Electronics | Online  | 4  (100) | 490_000   | category + channel
NULL    | Electronics | NULL    | 5  (101) | 680_000   | только category
NULL    | NULL        | Online  | 6  (110) | 810_000   | только channel
NULL    | NULL        | NULL    | 7  (111) | 1_250_000 | Grand Total

level_id = 7 (все биты = 1, все колонки агрегированы) - Grand Total. level_id = 0 (все биты = 0, все колонки реальны) - детальный уровень.

Это позволяет фильтровать конкретные уровни аналитики:

-- Только страновой уровень (country реален, остальные агрегированы)
WHERE GROUPING_ID(country, category, channel) = 3  -- (011) в двоичной

-- Только Grand Total
WHERE GROUPING_ID(country, category, channel) = 7  -- (111) в двоичной

-- Только детальный уровень (без агрегатов)
WHERE GROUPING_ID(country, category, channel) = 0  -- (000) в двоичной

В DataFrame API:

from pyspark.sql import functions as F

result = df.cube("country", "category", "channel") \
    .agg(
        F.sum("revenue").alias("revenue"),
        F.grouping_id("country", "category", "channel").alias("level_id")
    )

# Фильтрация по уровню
country_level = result.filter(F.col("level_id") == 3)
grand_total   = result.filter(F.col("level_id") == 7)

Практическое применение: замена NULL на читаемые метки

После получения результата технические NULL можно заменить на понятные строки:

result = df.rollup("country", "category") \
    .agg(
        F.sum("revenue").alias("revenue"),
        F.grouping_id("country", "category").alias("level_id")
    ) \
    .withColumn(
        "country_label",
        F.when(F.col("level_id").bitwiseAND(2) > 0, F.lit("All Countries"))
         .otherwise(F.col("country"))
    ) \
    .withColumn(
        "category_label",
        F.when(F.col("level_id").bitwiseAND(1) > 0, F.lit("All Categories"))
         .otherwise(F.col("category"))
    )

Для level_id бит 1 (двоичное) означает category агрегирована, бит 2 - country агрегирована. bitwiseAND позволяет проверить конкретный бит без необходимости перечислять все возможные значения level_id.


Под капотом: оператор Expand

Как Spark реализует многомерную агрегацию без многократного чтения данных? Ответ - оператор Expand.

df.rollup("country", "category").agg(F.sum("revenue")).explain(mode="formatted")
== Physical Plan ==
*(3) HashAggregate(keys=[country#10, category#11, spark_grouping_id#99], functions=[sum(revenue#12)])
+- Exchange hashpartitioning(country#10, category#11, spark_grouping_id#99, 200)
   +- *(2) HashAggregate(keys=[country#10, category#11, spark_grouping_id#99], functions=[partial_sum(revenue#12)])
      +- *(2) Expand [List(country#10, category#11, 0, revenue#12),
                      List(country#10, null, 1, revenue#12),
                      List(null, null, 3, revenue#12)],
               [country#10, category#11, spark_grouping_id#99, revenue#12]
         +- *(1) Scan parquet ...

Что делает Expand? Для каждой входной строки создаёт N строк - по одной для каждой комбинации агрегации. Для ROLLUP(country, category) - три строки из каждой:

  1. Оригинальная строка: (country, category, 0) - детальный уровень
  2. Агрегированная по category: (country, NULL, 1) - промежуточный итог
  3. Grand Total: (NULL, NULL, 3) - итого

spark_grouping_id добавляется автоматически и используется для правильной агрегации (группировка включает этот ID, чтобы (RU, NULL) для итога по России не смешивалась с (NULL, NULL) Grand Total).

Следствие для производительности: Expand увеличивает объём данных, поступающих в HashAggregate. Для ROLLUP с 3 измерениями - в 4 раза (4 комбинации). Для CUBE с 3 измерениями - в 8 раз. Это увеличивает:

  • Объём данных в частичной агрегации
  • Объём Shuffle (Exchange после HashAggregate partial)
  • Риск Data Skew - если отдельные ключи агрегации широко представлены

Сравнение с UNION ALL: при 4 разрезах UNION ALL читает данные 4 раза, но без Expand. GROUPING SETS читает данные 1 раз, но с Expand умножает объём на этапе агрегации. На практике GROUPING SETS быстрее, потому что I/O (чтение файлов) дороже памяти и CPU.


Физический план: что смотреть в explain

После понимания оператора Expand можно читать план целенаправленно:

sales = spark.read.parquet("s3://bucket/sales/")

result = sales \
    .filter(F.col("year") == 2024) \
    .cube("country", "category", "channel") \
    .agg(
        F.sum("revenue").alias("revenue"),
        F.avg("order_value").alias("avg_check"),
        F.count("*").alias("orders")
    )

result.explain(mode="formatted")

Что искать в плане:

1. Оператор Expand - должен быть ровно один. Проверьте число списков в Expand [List(...), List(...), ...]. Для CUBE(3 колонки) - 8 списков, для ROLLUP(3 колонки) - 4 списка.

2. Два HashAggregate - нормальная двухфазная агрегация (partial + final). Если вместо HashAggregate видите SortAggregate - это медленнее (см. предыдущий урок).

3. Один Exchange - после первого (partial) HashAggregate. Shuffle по составному ключу: (country, category, channel, spark_grouping_id).

4. PushedFilters в Scan - убедитесь, что year = 2024 попало в PushedFilters. Если фильтр не попал - Spark читает все данные за все годы.

Типичный план для cube("a", "b", "c"):

*(3) HashAggregate(keys=[a, b, c, gid], functions=[sum, avg, count])
+- Exchange hashpartitioning(a, b, c, gid, 200)
   +- *(2) HashAggregate(keys=[a, b, c, gid], functions=[partial_sum, partial_avg, partial_count])
      +- *(2) Expand [
             List(a, b, c, 0, revenue, order_value),  -- (a,b,c) детальный
             List(a, b, null, 1, revenue, order_value), -- (a,b)
             ...8 всего для CUBE...
         ], [a, b, c, gid, revenue, order_value]
         +- *(1) Filter (year = 2024)
            +- *(1) Scan parquet ... PushedFilters: [EqualTo(year,2024)]

Полный практический пример: e-commerce витрина

Разберём реальный кейс: построение аналитической витрины для e-commerce с четырьмя уровнями агрегации.

from pyspark.sql import functions as F
from pyspark.sql.window import Window

# Входные данные
sales = spark.read.parquet("s3://datalake/silver/sales/") \
    .filter(F.col("year") == 2024)

# Многомерная агрегация: страна, категория, канал продаж
cube_result = (
    sales.cube("country", "category", "channel")
    .agg(
        F.sum("revenue").alias("total_revenue"),
        F.sum("cost").alias("total_cost"),
        F.avg("order_value").alias("avg_check"),
        F.count("*").alias("order_count"),
        F.countDistinct("user_id").alias("unique_customers"),
        F.grouping_id("country", "category", "channel").alias("level_id")
    )
)

# Добавить читаемые метки вместо технических NULL
def label_aggregated(col_name: str, bit_position: int):
    """Заменить NULL агрегированного измерения на читаемую метку."""
    return F.when(
        F.col("level_id").bitwiseAND(bit_position) > 0,
        F.lit(f"All {col_name}s")
    ).otherwise(F.col(col_name.lower()))

vitrine = (
    cube_result
    .withColumn("country_label",  label_aggregated("Country",  4))  # бит 2 (позиция 2)
    .withColumn("category_label", label_aggregated("Category", 2))  # бит 1 (позиция 1)
    .withColumn("channel_label",  label_aggregated("Channel",  1))  # бит 0 (позиция 0)
    .withColumn(
        "level_name",
        F.when(F.col("level_id") == 7, F.lit("Grand Total"))
         .when(F.col("level_id") == 3, F.lit("Country Total"))
         .when(F.col("level_id") == 1, F.lit("Country+Category"))
         .when(F.col("level_id") == 0, F.lit("Detail"))
         .otherwise(F.concat(F.lit("level_"), F.col("level_id").cast("string")))
    )
    .withColumn(
        "gross_margin_pct",
        F.round(
            (F.col("total_revenue") - F.col("total_cost")) / F.col("total_revenue") * 100,
            2
        )
    )
    .drop("level_id", "country", "category", "channel")
    .withColumnRenamed("country_label",  "country")
    .withColumnRenamed("category_label", "category")
    .withColumnRenamed("channel_label",  "channel")
)

# Записать витрину в Gold-слой
vitrine.write \
    .mode("overwrite") \
    .partitionBy("level_name") \
    .parquet("s3://datalake/gold/sales_cube/year=2024/")

Обратите внимание на партиционирование по level_name: BI-инструмент, запрашивающий только Country Total, прочитает только соответствующую партицию - без сканирования детальных данных.


Временные ряды: ROLLUP для иерархии дат

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

from pyspark.sql import functions as F

sales = spark.read.parquet("s3://datalake/silver/sales/")

# Добавить временные измерения
enriched = sales.withColumn("year",    F.year("order_date")) \
                .withColumn("quarter", F.quarter("order_date")) \
                .withColumn("month",   F.month("order_date"))

# ROLLUP по временной иерархии
time_report = enriched.rollup("year", "quarter", "month") \
    .agg(
        F.sum("revenue").alias("revenue"),
        F.sum("orders").alias("orders"),
        F.avg("order_value").alias("avg_check"),
        F.grouping_id("year", "quarter", "month").alias("time_level")
    ) \
    .withColumn(
        "period_label",
        F.when(F.col("time_level") == 7, F.lit("Grand Total"))
         .when(F.col("time_level") == 3, F.concat(F.lit("Year "), F.col("year").cast("string")))
         .when(F.col("time_level") == 1, F.concat(
             F.lit("Q"), F.col("quarter").cast("string"),
             F.lit(" "), F.col("year").cast("string")
         ))
         .when(F.col("time_level") == 0, F.concat(
             F.col("year").cast("string"), F.lit("-"),
             F.lpad(F.col("month").cast("string"), 2, "0")
         ))
         .otherwise(F.lit("Unknown"))
    ) \
    .orderBy(
        F.col("time_level").desc(),  # Grand Total → Year → Quarter → Month
        F.col("year").asc_nulls_last(),
        F.col("quarter").asc_nulls_last(),
        F.col("month").asc_nulls_last()
    )

Результат:

year | quarter | month | time_level | period_label | revenue
NULL | NULL    | NULL  | 7          | Grand Total  | 12_500_000
2023 | NULL    | NULL  | 3          | Year 2023    |  4_800_000
2023 | 1       | NULL  | 1          | Q1 2023      |  1_200_000
2023 | 1       | 1     | 0          | 2023-01      |    400_000
2023 | 1       | 2     | 0          | 2023-02      |    380_000
2023 | 1       | 3     | 0          | 2023-03      |    420_000
...

GROUPING SETS в Spark SQL: продвинутый синтаксис

Spark SQL поддерживает полный стандарт SQL:2003 для расширенных группировок. Несколько полезных расширений:

Вложенные GROUPING SETS:

-- Агрегация по двум независимым иерархиям в одном запросе
SELECT country, category, year, month, SUM(revenue)
FROM sales
GROUP BY GROUPING SETS (
    (country, category),  -- географо-продуктовый разрез
    (year, month)         -- временной разрез
)
-- Ни один разрез не смешивает измерения из разных групп

Комбинирование ROLLUP и CUBE:

-- Иерархия по времени + полный куб по продукту
SELECT year, month, category, brand, SUM(revenue)
FROM sales
GROUP BY ROLLUP(year, month), CUBE(category, brand)
-- Эквивалентно: комбинация из (ROLLUP) × (CUBE) комбинаций

FILTER внутри агрегатных функций:

-- Разные условия для разных агрегатов в одном GROUP BY
SELECT
    country,
    SUM(revenue) AS total_revenue,
    SUM(revenue) FILTER (WHERE channel = 'Online') AS online_revenue,
    SUM(revenue) FILTER (WHERE channel = 'Offline') AS offline_revenue
FROM sales
GROUP BY ROLLUP(country)

В DataFrame API FILTER внутри agg:

df.groupBy("country").agg(
    F.sum("revenue").alias("total_revenue"),
    F.sum(
        F.when(F.col("channel") == "Online", F.col("revenue")).otherwise(F.lit(0))
    ).alias("online_revenue"),
    F.sum(
        F.when(F.col("channel") == "Offline", F.col("revenue")).otherwise(F.lit(0))
    ).alias("offline_revenue")
)

NULL-семантика при реальных NULL в данных

Разберём кейс, который часто создаёт проблемы: в исходных данных есть строки с NULL в группировочных колонках.

# Исходные данные: часть строк без страны (неизвестный источник)
raw_data = spark.createDataFrame([
    ("RU",   "Electronics", 1000.0),
    ("RU",   "Clothing",     500.0),
    (None,   "Electronics",  300.0),   # страна неизвестна
    ("DE",   "Electronics",  800.0),
], ["country", "category", "revenue"])

# ROLLUP(country, category)
result = raw_data.rollup("country", "category") \
    .agg(
        F.sum("revenue").alias("revenue"),
        F.grouping("country").alias("country_agg"),
        F.grouping("category").alias("category_agg"),
        F.grouping_id("country", "category").alias("level_id")
    )
result.orderBy("country", "category").show()
country | category    | revenue | country_agg | category_agg | level_id
NULL    | Electronics | 300.0   | 0           | 0            | 0  ← реальный NULL страны
NULL    | NULL        | 300.0   | 0           | 1            | 1  ← итого для NULL-страны
NULL    | NULL        | 2600.0  | 1           | 1            | 3  ← Grand Total
RU      | Clothing    | 500.0   | 0           | 0            | 0
RU      | Electronics | 1000.0  | 0           | 0            | 0
RU      | NULL        | 1500.0  | 0           | 1            | 1
DE      | Electronics | 800.0   | 0           | 0            | 0
DE      | NULL        | 800.0   | 0           | 1            | 1

Ключевые наблюдения:

  • Строка (NULL, Electronics, 300.0, country_agg=0) - это реальная запись с неизвестной страной. country_agg=0 означает: NULL здесь не от ROLLUP, а реальный NULL.
  • Строка (NULL, NULL, 300.0, country_agg=0) - промежуточный итог по "стране NULL". Тоже country_agg=0.
  • Строка (NULL, NULL, 2600.0, country_agg=1) - Grand Total. country_agg=1 отличает его от двух предыдущих.

Без GROUPING() функции три строки с country=NULL были бы неразличимы визуально.


Оптимизация производительности

Фильтруйте до агрегации

Predicate Pushdown работает до Expand - фильтры применяются к исходным данным, а не к умноженным Expand строкам:

# ХОРОШО: filter до cube - меньше данных входит в Expand
sales.filter(F.col("year") == 2024) \
     .cube("country", "category") \
     .agg(F.sum("revenue"))

# ПЛОХО (логически идентично, но убедитесь что Spark оптимизирует):
sales.cube("country", "category") \
     .agg(F.sum("revenue")) \
     .filter(F.col("year") == 2024)   # фильтр ПОСЛЕ куба

Всегда проверяйте через explain(), что PushedFilters в Scan содержит ваши условия.

Кешируйте при многократном использовании

Если вы строите несколько аналитических агрегатов из одного DataFrame:

# Кешировать исходный DataFrame перед множественными аналитическими запросами
base = sales.filter(F.col("year") == 2024).cache()
base.count()   # материализовать кеш

geo_cube    = base.cube("country", "region").agg(...)
time_rollup = base.rollup("month", "week").agg(...)
product_gs  = base.groupBy(F.grouping_sets(["category"], ["brand"], [])).agg(...)

Без кеша каждый cube()/rollup() будет читать данные с диска заново.

Ограничивайте кардинальность в CUBE

Если одно из измерений имеет много уникальных значений (например, user_id или product_sku) - не включайте его в CUBE. Используйте GROUPING SETS и явно исключите высококардинальные комбинации:

# ПЛОХО: CUBE с тысячами уникальных user_id
df.cube("country", "category", "user_id").agg(F.sum("revenue"))
# Результат: тысячи × тысячи строк, огромный Shuffle

# ХОРОШО: только нужные комбинации без user_id
df.groupBy(F.grouping_sets(
    [],
    ["country"],
    ["category"],
    ["country", "category"]
)).agg(F.sum("revenue"))

Настройте число партиций Shuffle

После Expand объём данных умножается. Дефолтные 200 партиций могут оказаться много (пустые партиции) или мало (большие партиции):

# Для CUBE с небольшими данными (< 1 GB после фильтрации)
spark.conf.set("spark.sql.shuffle.partitions", "20")

# Для больших аналитических таблиц
spark.conf.set("spark.sql.shuffle.partitions", "400")

# С AQE - Spark сам скорректирует
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")

Практика: переписать UNION ALL на GROUPING SETS

Задача: в команде есть legacy-код с 5 блоками UNION ALL. Переписать на GROUPING SETS и сравнить планы.

Было:

# Legacy: 5 отдельных groupBy + UNION ALL
q1 = df.agg(F.sum("revenue").alias("rev")) \
       .withColumn("country", F.lit(None)) \
       .withColumn("category", F.lit(None)) \
       .withColumn("channel", F.lit(None))

q2 = df.groupBy("country").agg(F.sum("revenue").alias("rev")) \
       .withColumn("category", F.lit(None)) \
       .withColumn("channel", F.lit(None))

q3 = df.groupBy("category").agg(F.sum("revenue").alias("rev")) \
       .withColumn("country", F.lit(None)) \
       .withColumn("channel", F.lit(None))

q4 = df.groupBy("country", "category").agg(F.sum("revenue").alias("rev")) \
       .withColumn("channel", F.lit(None))

q5 = df.groupBy("country", "category", "channel").agg(F.sum("revenue").alias("rev"))

legacy_result = q1.unionByName(q2).unionByName(q3).unionByName(q4).unionByName(q5)

Стало:

-- Spark SQL
SELECT
    country,
    category,
    channel,
    SUM(revenue) AS rev,
    GROUPING_ID(country, category, channel) AS level_id
FROM sales
GROUP BY GROUPING SETS (
    (),
    (country),
    (category),
    (country, category),
    (country, category, channel)
)
# DataFrame API
from pyspark.sql.functions import grouping_sets

new_result = df.groupBy(
    grouping_sets(
        [],
        ["country"],
        ["category"],
        ["country", "category"],
        ["country", "category", "channel"]
    )
).agg(
    F.sum("revenue").alias("rev"),
    F.grouping_id("country", "category", "channel").alias("level_id")
)

Сравните планы:

# Legacy plan - 5 scan + 5 shuffle + union
legacy_result.explain(mode="formatted")

# Новый план - 1 scan + 1 expand + 1 shuffle
new_result.explain(mode="formatted")

В плане legacy вы увидите 5 блоков Scan → HashAggregate → Exchange соединённых через Union. В новом плане - один блок Scan → Expand → HashAggregate → Exchange → HashAggregate.


Антипаттерны многомерной агрегации

1. CUBE с высококардинальными измерениями.

CUBE(country, category, user_id) при 50 000 уникальных user_id даст 50 000 × количество_стран × количество_категорий строк только на уровне (country, category, user_id). Плюс все остальные 7 комбинаций. OOM на Driver при materialize.

2. CUBE без предварительной фильтрации.

Если в таблице данные за 5 лет, а нужны только данные за последний квартал - фильтруйте сразу. Без фильтра Expand умножает все строки.

3. Игнорирование GROUPING_ID при последующем использовании.

Результат ROLLUP/CUBE без GROUPING_ID или GROUPING() - таблица с неразличимыми NULL. Скармливать такой результат BI-инструменту напрямую - создавать путаницу для пользователей.

4. Слишком много агрегатных функций.

Каждая дополнительная агрегатная функция увеличивает размер строки, которую Expand умножает. 10 sum()/avg()/count() в одном CUBE - это нормально. 50+ функций с collect_list и percentile - переосмыслите подход.

5. Использование CUBE вместо ROLLUP для иерархических данных.

CUBE для (year, quarter, month) создаёт комбинации (NULL, quarter, NULL) - итог по кварталу без привязки к году. Это аналитически бессмысленно. Используйте ROLLUP.


Итоги: чек-лист аналитических агрегаций

Выбор оператора:

  • Иерархические данные (год → месяц → день) → ROLLUP
  • Все комбинации равноправных измерений → CUBE (не более 4–5 колонок)
  • Конкретный список разрезов → GROUPING SETS
  • Несколько несвязанных разрезов → GROUPING SETS вместо UNION ALL

Идентификация уровней:

  • Всегда добавлять GROUPING_ID() в результат
  • Использовать GROUPING() для замены NULL на читаемые метки
  • Партиционировать Gold-таблицу по level_name для эффективного чтения

Производительность:

  • Фильтровать до агрегации, проверять PushedFilters в explain()
  • Кешировать исходный DataFrame при множественных запросах
  • Не превышать 4–5 измерений в CUBE
  • Настраивать spark.sql.shuffle.partitions соразмерно объёму данных
  • Включить AQE для автоматического coalescing пустых партиций