大規模な PySpark

PySpark 入門

Benjamin Schmidt

Data Engineer

スケールの活用

  • Pysparkはギガバイトやテラバイトのデータを効率的に処理できます
  • PySpark を使用すると、速度と効率的な処理が目標です
  • PySpark の実行を理解すると、さらに効率化できる
  • ブロードキャストを使用してクラスター全体を管理する
joined_df = large_df.join(broadcast(small_df), 
                          on="key_column", how="inner")
joined_df.show()
PySpark 入門

実行計画

# Using explain() to view the execution plan
df.filter(df.Age > 40).select("Name").explain()
== Physical Plan ==
*(1) Filter (isnotnull(Age) AND (Age > 30))
+- Scan ExistingRDD[Name:String, Age:Int]
1 https://spark.apache.org/docs/latest/api/python/reference/pyspark.sql/api/pyspark.sql.DataFrame.explain.html
PySpark 入門

DataFrame のキャッシュと永続化

  • キャッシュ: より高速なアクセスのため、小規模なデータセットのデータをメモリに保存します
  • 永続化: 異なるストレージレベルにデータを保存します
df = spark.read.csv("large_dataset.csv", header=True, inferSchema=True)

# Cache the DataFrame
df.cache()

# Perform multiple operations on the cached DataFrame df.filter(df["column1"] > 50).show() df.groupBy("column2").count().show()
PySpark 入門

異なるストレージレベルでDataFrameを保存する

# Persist the DataFrame with storage level
from pyspark import StorageLevel

df.persist(StorageLevel.MEMORY_AND_DISK)

# Perform transformations result = df.groupBy("column3").agg({"column4": "sum"}) result.show() # Unpersist after use df.unpersist()
PySpark 入門

PySparkの最適化

  • 小さなサブセクション: 使用されるデータが多いほど、処理は遅くなります。 選択性の高い手法のため、map()よりgroupby()のようなツールを選んでください
  • ブロードキャスト結合: ブロードキャストは、小規模なデータセットでもすべての計算リソースを使用します
  • 繰り返しの操作を避ける 同じデータに対する繰り返しの操作は、何の利点もなく時間とコンピュートを消費します
PySpark 入門

練習しましょう!

PySpark 入門

Preparing Video For Download...