Агрегации в PySpark

Введение в PySpark

Benjamin Schmidt

Data Engineer

Обзор SQL-агрегаций в PySpark

  • Стандартные SQL-агрегации работают с spark.sql()
    # SQL aggregation query
    spark.sql("""
      SELECT Department, SUM(Salary) AS Total_Salary, AVG(Salary) AS Average_Salary
      FROM employees
      GROUP BY Department
    """).show()
    
Введение в PySpark

Совместное использование DataFrame и SQL

# Filter salaries over 3000
filtered_df = df.filter(df.Salary > 3000)

# Register filtered DataFrame as a view
filtered_df.createOrReplaceTempView("filtered_employees")

# Aggregate using SQL on the filtered view spark.sql(""" SELECT Department, COUNT(*) AS Employee_Count FROM filtered_employees GROUP BY Department """).show()
Введение в PySpark

Работа с типами данных при агрегации

# Example of type casting
data = [("HR", "3000"), ("IT", "4000"), ("Finance", "3500")]
columns = ["Department", "Salary"]
df = spark.createDataFrame(data, schema=columns)

# Convert Salary column to integer df = df.withColumn("Salary", df["Salary"].cast("int")) # Perform aggregation df.groupBy("Department").sum("Salary").show()
Введение в PySpark

RDD для агрегаций

# Example of aggregation with RDDs
rdd = df.rdd.map(lambda row: (row["Department"], row["Salary"]))

rdd_aggregated = rdd.reduceByKey(lambda x, y: x + y)
print(rdd_aggregated.collect())
Введение в PySpark

Лучшие практики агрегаций в PySpark

  • Фильтруйте заранее: уменьшайте объём данных до агрегации
  • Следите за типами: данные должны быть чистыми и корректно типизированными
  • Избегайте операций над всем набором данных: минимизируйте использование groupBy()
  • Выбирайте подходящий инструмент: предпочитайте DataFrame — они лучше оптимизированы
  • Контролируйте производительность: используйте explain() для анализа плана выполнения
Введение в PySpark

Основные выводы

  • SQL-агрегации в PySpark: функции SUM() и AVERAGE() для обобщения данных
  • DataFrame и SQL: сочетание обоих подходов для гибкой обработки данных
  • Работа с типами данных: решение проблем несовместимости типов при агрегации
  • RDD и DataFrame: понимание компромиссов и выбор подходящего инструмента
Введение в PySpark

Давайте потренируемся!

Введение в PySpark

Preparing Video For Download...