Группировка данных в PySpark
groupBy() по одной и нескольким колонкам, модуль F (pyspark.sql.functions), агрегирующие функции count/sum/avg/min/max/countDistinct, метод agg() для множественных агрегаций.
Группировка с groupBy()¶
В PySpark данные группируются методом groupBy(). Он группирует строки по значениям одной или нескольких колонок.
# Группировка по одной колонке
grouped_df = df.groupBy("department")
# Группировка по нескольким колонкам
grouped_df = df.groupBy("department", "location")
После группировки выполняются агрегации (подсчёт, суммирование и др.) по каждой группе. Используйте alias() для переименования агрегированной колонки.
Зачем использовать F для функций?¶
В PySpark большинство агрегирующих функций доступны в модуле pyspark.sql.functions. Принято импортировать этот модуль как F по двум причинам:
- Ясность: отличает функции PySpark (например,
F.sum) от встроенных функций Python (например,sum). - Читаемость: использование
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 легко справляется как с простыми подсчётами, так и со сложными многопоказательными агрегациями на больших наборах данных.