Стійкі розподілені набори даних у PySpark

Вступ до PySpark

Benjamin Schmidt

Data Engineer

What is parallelization in 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 проти DataFrame

DataFrame

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

RDD

  • Низький рівень: гнучкіші, але для складних операцій потрібно більше коду
  • Типобезпека: зберігають типи даних, але не мають оптимізацій DataFrame
  • Без схеми: складніше працювати зі структурованими або реляційними даними, як у SQL
  • Масштабування на великі обсяги
  • Дуже багатослівні порівняно з DataFrame і слабкі для аналітики
Вступ до PySpark

Корисні функції та методи

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

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

Вступ до PySpark

Preparing Video For Download...