PySparkの集計

PySpark 入門

Benjamin Schmidt

Data Engineer

PySpark SQL 集計の概要

  • 一般的な 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 入門

重要ポイント

  • PySpark SQL 集計: データを要約するためのSUM()AVERAGE()のような機能
  • DataFrame と SQL: 柔軟なデータ操作のために両方のアプローチを組み合わせる
  • データ型の処理: 集計処理中の型不一致の問題への対処
  • RDD と DataFrame の比較: トレードオフを理解し、適切なツールを選ぶこと
PySpark 入門

練習しましょう!

PySpark 入門

Preparing Video For Download...