Группировка данных в PySpark

groupBy() по одной и нескольким колонкам, модуль F (pyspark.sql.functions), агрегирующие функции count/sum/avg/min/max/countDistinct, метод agg() для множественных агрегаций.

core

Группировка с groupBy()

В PySpark данные группируются методом groupBy(). Он группирует строки по значениям одной или нескольких колонок.

# Группировка по одной колонке
grouped_df = df.groupBy("department")

# Группировка по нескольким колонкам
grouped_df = df.groupBy("department", "location")

После группировки выполняются агрегации (подсчёт, суммирование и др.) по каждой группе. Используйте alias() для переименования агрегированной колонки.

Зачем использовать F для функций?

В PySpark большинство агрегирующих функций доступны в модуле pyspark.sql.functions. Принято импортировать этот модуль как F по двум причинам:

  1. Ясность: отличает функции PySpark (например, F.sum) от встроенных функций Python (например, sum).
  2. Читаемость: использование F делает код лаконичным и понятным.
from pyspark.sql import functions as F

# Одна агрегация
grouped_sum = df.groupBy("department").sum("salary")

# Через F с переименованием агрегированной колонки
aggregated_df = df.groupBy("department").agg(F.sum("salary").alias("total_salary"))

# Несколько агрегаций
aggregated_df = df.groupBy("department").agg(
    F.sum("salary").alias("total_salary"),
    F.avg("salary").alias("average_salary"),
    F.count("*").alias("employee_count")
)

Основные агрегирующие функции

1. count() — подсчёт строк в группе

# Подсчёт строк по отделам
grouped_count = df.groupBy("department").count()

2. sum() — сумма числовой колонки

# Сумма зарплат по отделам
grouped_sum = df.groupBy("department").sum("salary")

3. avg() — среднее числовой колонки

# Средняя зарплата по отделам
grouped_avg = df.groupBy("department").avg("salary")

4. min() и max() — минимум или максимум в группе

# Минимальная и максимальная зарплата по отделам
grouped_min_max = df.groupBy("department").min("salary").max("salary")

5. countDistinct() — подсчёт уникальных значений

# Подсчёт уникальных должностей по отделам
distinct_count = df.groupBy("department").agg(F.countDistinct("job_role").alias("unique_roles"))

Несколько агрегаций через agg()

Для одновременного выполнения нескольких агрегаций используйте метод agg():

# Несколько агрегаций
aggregated_df = df.groupBy("department").agg(
    F.avg("salary").alias("average_salary"),
    F.sum("bonus").alias("total_bonus")
)

Полный пример

from pyspark.sql import SparkSession
from pyspark.sql.functions import avg, sum, count

# Инициализация Spark Session
spark = SparkSession.builder \
    .appName("PySpark Grouping Example") \
    .getOrCreate()

# Тестовый DataFrame
data = [
    ("Alice", "HR", 5000),
    ("Bob", "IT", 6000),
    ("Charlie", "Finance", 7000),
    ("David", "IT", 6000),
    ("Eve", "HR", 5500),
    ("Frank", "Finance", 8000),
]
columns = ["name", "department", "salary"]

df = spark.createDataFrame(data, columns)

# Показать исходные данные
print("Original Data:")
df.show()

# Группировка по отделу и вычисление агрегатов
print("Group by Department - Count:")
df.groupBy("department").count().show()

print("Group by Department - Sum of Salaries:")
df.groupBy("department").sum("salary").show()

print("Group by Department - Average Salary:")
df.groupBy("department").agg(avg("salary")).show()

print("Group by Department - Multiple Aggregates:")
df.groupBy("department").agg(
    count("name").alias("employee_count"),
    sum("salary").alias("total_salary"),
    avg("salary").alias("average_salary")
).show()

# Остановка Spark Session
spark.stop()

Итог

Группировка в PySpark — мощный способ агрегировать данные, аналогичный SQL GROUP BY. Модуль F обеспечивает доступ к агрегирующим функциям, сохраняя ясность и лаконичность кода. PySpark легко справляется как с простыми подсчётами, так и со сложными многопоказательными агрегациями на больших наборах данных.