Аналитические запросы: GROUPING SETS, ROLLUP, CUBE без UDF
GROUPING SETS, ROLLUP, CUBE - многомерная агрегация за один проход. GROUPING_ID для идентификации уровней, физический план Expand, сравнение с UNION ALL.
Проблема: 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) - три строки из каждой:
- Оригинальная строка:
(country, category, 0)- детальный уровень - Агрегированная по category:
(country, NULL, 1)- промежуточный итог - 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 пустых партиций