Устойчивые распределённые наборы данных в PySpark

Введение в PySpark

Benjamin Schmidt

Data Engineer

Что такое параллелизация в PySpark?

  • Автоматическое распараллеливание данных и вычислений по узлам кластера
  • Распределённая обработка больших наборов данных на нескольких узлах
  • Рабочие узлы обрабатывают данные параллельно и объединяют результаты
  • Высокая скорость при больших объёмах (гигабайты и терабайты)

Параллелизация

Введение в PySpark

Знакомство с RDD

RDD (Resilient Distributed Datasets):

  • Распределённые коллекции данных в кластере с автоматическим восстановлением при сбоях узлов
  • Подходят для работы с большими данными
  • Неизменяемы; поддерживают преобразования — map(), filter() — и действия — collect(), paralelize() — для получения результатов или создания RDD
Введение в PySpark

Создание RDD

# Initialize a Spark session
from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("RDDExample").getOrCreate()

# Create a DataFrame from a csv census_df = spark.read.csv("/census.csv")
# Convert DataFrame to RDD census_rdd = census_df.rdd
# Show the RDD's contents using collect() census_rdd.collect()
Введение в PySpark

Метод collect()

# Collect the entire DataFrame into a local Python list of Row objects
data_collected = df.collect()

# Print the collected data
for row in data_collected:
    print(row)
```    
Введение в PySpark

RDD и DataFrames

DataFrames

  • Высокоуровневые: оптимизированы для удобства работы
  • SQL-подобные операции: поддерживают запросы в стиле SQL и сложные операции с минимальным кодом
  • Схема данных: содержат столбцы и типы, как таблица SQL

RDD

  • Низкоуровневые: гибкие, но требуют больше кода для сложных операций
  • Типобезопасность: сохраняют типы данных, но лишены оптимизаций DataFrames
  • Без схемы: сложнее работать со структурированными и реляционными данными
  • Хорошо масштабируются
  • Очень многословны по сравнению с DataFrames и слабо подходят для аналитики
Введение в PySpark

Полезные функции и методы

  • map(): применяет функции (в том числе пользовательские, например лямбда-функции) к набору данных: rdd.map(map_function)
  • collect(): собирает данные со всего кластера: rdd.collect()
Введение в PySpark

Давайте потренируемся!

Введение в PySpark

Preparing Video For Download...