Подробнее о Spark DataFrames

Введение в PySpark

Benjamin Schmidt

Data Engineer

Создание DataFrames из различных источников данных

  • CSV-файлы: стандартный формат для структурированных данных с разделителями
  • JSON-файлы: полуструктурированный иерархический формат
  • Parquet-файлы: оптимизированы для хранения и запросов, широко используются в инженерии данных
  • Пример:
    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

  • Схему можно задать вручную для большего контроля — удобно при фиксированных структурах данных

Схема в масштабе

Введение в PySpark

Типы данных в PySpark DataFrames

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

Синтаксис типов данных для PySpark DataFrames

# Import the necessary types as classes
from pyspark.sql.types import (StructType,
                            StructField, IntegerType,
                            StringType, ArrayType)

# Construct the schema
schema = StructType([
    StructField("id", IntegerType(), True),
    StructField("name", StringType(), True),
    StructField("scores", ArrayType(IntegerType()), True)
])

# Set the schema
df = spark.createDataFrame(data, schema=schema)
Введение в PySpark

Операции с DataFrame — выборка и фильтрация

  • .select() — выбор нужных столбцов
  • .filter() или .where() — фильтрация строк по условию
  • .sort() — сортировка по набору столбцов
# Select and show only the name and age columns
df.select("name", "age").show()
# Filter on age > 30
df.filter(df["age"] > 30).show()
# Use Where to filter match a specific value
df.where(df["age"] == 30).show()
# Use Sort to sort on age
df.sort("age", ascending=False).show()
Введение в PySpark

Сортировка и удаление пропущенных значений

  • Сортировка данных с помощью .sort() или .orderBy()
  • Удаление строк с пустыми значениями через na.drop()
# Sort using the age column
df.sort("age", ascending=False).show()

# Drop missing values
df.na.drop().show()

Введение в PySpark

Шпаргалка

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

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

Введение в PySpark

Preparing Video For Download...