Більше про DataFrame у Spark

Вступ до PySpark

Benjamin Schmidt

Data Engineer

Створення DataFrame з різних джерел даних

  • Файли CSV: поширені для структурованих, розділених даних
  • Файли JSON: напівструктурований, ієрархічний формат
  • Файли Parquet: оптимізовані для зберігання й запитів; часто у data engineering
  • Приклад:
    spark.read.csv("path/to/file.csv")
    
  • Приклад:
    spark.read.json("path/to/file.json")
    
  • Приклад:
    spark.read.parquet("path/to/file.parquet")
    
1 https://spark.apache.org/docs/latest/api/python/reference/pyspark.pandas/api/pyspark.pandas.read_csv
Вступ до PySpark

Визначення схеми: автоматично й вручну

  • Spark може виводити схеми з даних за допомогою inferSchema=True

  • Визначайте схему вручну для кращого контролю — корисно для фіксованих структур

Schema at Scale

Вступ до PySpark

Типи даних у DataFrame PySpark

  • IntegerType: цілі числа
    • Напр., 1, 3478, -1890456
  • LongType: більші цілі числа
    • Напр., 8-байтові зі знаком, 922334775806
  • FloatType і DoubleType: числа з плаваючою комою
    • Напр., 3.14159
  • StringType: текстові рядки
    • Напр., "This is an example of a string."
  • ...
Вступ до PySpark

Синтаксис типів даних для DataFrame у PySpark

# Імпортуйте необхідні типи як класи
from pyspark.sql.types import (StructType,
                            StructField, IntegerType,
                            StringType, ArrayType)

# Створіть схему
schema = StructType([
    StructField("id", IntegerType(), True),
    StructField("name", StringType(), True),
    StructField("scores", ArrayType(IntegerType()), True)
])

# Задайте схему
df = spark.createDataFrame(data, schema=schema)
Вступ до PySpark

Операції з DataFrame: вибір і фільтрація

  • Використовуйте .select() для вибору колонок
  • Використовуйте .filter() або .where() для фільтрації рядків за умовами
  • Використовуйте .sort() для сортування за кількома колонками
# Виберіть і покажіть лише колонки name та age
df.select("name", "age").show()
# Фільтруйте за age > 30
df.filter(df["age"] > 30).show()
# Використайте where для точного значення
df.where(df["age"] == 30).show()
# Сортуйте за age
df.sort("age", ascending=False).show()
Вступ до PySpark

Сортування і видалення пропусків

  • Упорядковуйте дані за допомогою .sort() або .orderBy()
  • Використовуйте na.drop() для видалення рядків із null-значеннями
# Сортування за колонкою age
df.sort("age", ascending=False).show()

# Видаліть пропущені значення
df.na.drop().show()

Вступ до PySpark

Шпаргалка

  • spark.read_json(): завантаження даних з JSON
  • spark.read.schema(): явне визначення схем
  • .na.drop(): видалення рядків із пропусками
  • .select(), .filter(), .sort(), .orderBy(): базові функції обробки даних
Вступ до PySpark

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

Вступ до PySpark

Preparing Video For Download...