Agregacje w PySpark

Wprowadzenie do PySpark

Benjamin Schmidt

Data Engineer

Przegląd agregacji SQL w PySpark

  • Typowe agregacje SQL działają z 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()
    
Wprowadzenie do PySpark

Łączenie operacji DataFrame i 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()
Wprowadzenie do PySpark

Obsługa typów danych w agregacjach

# 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()
Wprowadzenie do PySpark

RDD w agregacjach

# 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())
Wprowadzenie do PySpark

Najlepsze praktyki agregacji w PySpark

  • Wczesne filtrowanie: Zmniejszenie rozmiaru danych przed agregacją
  • Obsługa typów danych: Zapewnienie czystości i poprawności typów
  • Unikanie operacji na całym zbiorze: Minimalizowanie operacji jak groupBy()
  • Wybór właściwego interfejsu: Preferowanie DataFrames ze względu na optymalizacje
  • Monitorowanie wydajności: Użycie explain() do inspekcji i optymalizacji planu wykonania
Wprowadzenie do PySpark

Najważniejsze wnioski

  • Agregacje SQL w PySpark: Funkcje takie jak SUM() i AVERAGE() do podsumowywania danych
  • DataFrames i SQL: Łączenie obu podejść dla elastycznej manipulacji danymi
  • Obsługa typów danych: Rozwiązywanie problemów z niezgodnością typów podczas agregacji
  • RDD a DataFrames: Zrozumienie kompromisów i wybór właściwego narzędzia
Wprowadzenie do PySpark

Czas na ćwiczenia!

Wprowadzenie do PySpark

Preparing Video For Download...