PySpark у масштабі

Вступ до PySpark

Benjamin Schmidt

Data Engineer

Використання масштабу

  • PySpark ефективний для гігабайтів і терабайтів даних
  • У PySpark мета — швидкість і ефективна обробка
  • Розуміння виконання PySpark дає ще більше ефективності
  • Використовуйте broadcast для керування всім кластером
joined_df = large_df.join(broadcast(small_df), 
                          on="key_column", how="inner")
joined_df.show()
Вступ до PySpark

Плани виконання

# Використання explain() для перегляду плану виконання
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)

# Кешуйте DataFrame
df.cache()

# Виконайте кілька операцій на кешованому DataFrame df.filter(df["column1"] > 50).show() df.groupBy("column2").count().show()
Вступ до PySpark

Персистування DataFrame на різних рівнях сховища

# Персистування DataFrame з рівнем сховища
from pyspark import StorageLevel

df.persist(StorageLevel.MEMORY_AND_DISK)

# Виконайте перетворення result = df.groupBy("column3").agg({"column4": "sum"}) result.show() # Відв'яжіть після використання df.unpersist()
Вступ до PySpark

Оптимізація PySpark

  • Невеликі підмножини: що більше даних обробляється, то повільніша операція; віддавайте перевагу map() над groupby() завдяки вибірковості методів
  • Об'єднання з broadcast: broadcast задіє весь обчислювальний ресурс, навіть для менших наборів
  • Уникайте повторних дій: повторне виконання над тими самими даними марнує час і ресурси без користі
Вступ до PySpark

Давайте потренуємось!

Вступ до PySpark

Preparing Video For Download...