Введение в PySpark

Что такое PySpark, зачем он нужен и чем отличается от pandas. SparkSession, DataFrame, ленивые вычисления, трансформации и actions.

core

⚡ Введение в PySpark

Добро пожаловать в PySpark — продукт союза Python и Apache Spark! Python предоставляет языковую основу, а Apache Spark — мощный движок с открытым исходным кодом, созданный для обработки огромных объёмов данных на множестве машин с высокой скоростью.

PySpark vs pandas: если pandas — это езда на велосипеде, то PySpark — это прыжок в космический корабль, когда речь идёт о гигабайтах или терабайтах данных.

Установка PySpark 🛠️

Чтобы установить PySpark локально:

pip install pyspark

Знакомьтесь: SparkSession

SparkSession — это точка входа во всё, что связано со Spark.

from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("PySpark 101") \
    .getOrCreate()

print("Spark's in the house! 🔥")

DataFrame: как pandas, но на стероидах 💪

DataFrame в PySpark работает с наборами данных куда большего размера, чем pandas. Пример использования:

df = spark.read.csv("/path/to/your/fancy_file.csv", header=True, inferSchema=True)

df.show(5)

DataFrame ленив... но в хорошем смысле 🛌

PySpark использует ленивые вычисления (lazy evaluation). Операции вроде фильтрации не выполняются немедленно:

df_filtered = df.filter(df['age'] > 30)

Реальное вычисление запускается только при вызове действия (action):

df_filtered.show()

RDD: предки DataFrame 🧙‍♂️

RDD (Resilient Distributed Datasets) появились раньше DataFrame. Придерживайтесь DataFrame, если только вам не нужно ручное управление данными.

Трансформации и actions: инь и ян PySpark

Трансформации — ленивые операции, которые создают список инструкций:

  • .filter(), .select(), .groupBy()

Actions запускают фактическое выполнение:

  • .show(), .count(), .collect()

Пример: всё вместе 🎉

df = spark.read.csv("/path/to/employee_data.csv", header=True, inferSchema=True)

df.show(5)

adults = df.filter(df['age'] > 30)

department_count = adults.groupBy("department").count()

department_count.show()

Почему PySpark?

  • Скорость: параллельная обработка на множестве машин
  • Масштаб: работает с наборами данных от 10 МБ до 10 ТБ
  • Мощность: SQL-запросы, машинное обучение и аналитика в реальном времени