Фильтрация данных в PySpark

Методы filter() и where(), логические операторы AND/OR/NOT, isin(), startswith(), endswith(), like(), rlike(), isNull(), isNotNull(), SQL-запросы через createOrReplaceTempView.

core

1. Базовая фильтрация: filter() и where()

Основные методы фильтрации в PySpark — filter() и where(). Они функционально идентичны и принимают как выражения, так и условия для колонок.

# Пример базовой фильтрации
filtered_df = df.filter(df["age"] > 30)
# или
filtered_df = df.where("age > 30")

Для SQL-разработчиков метод .where() интуитивно понятен — он повторяет SQL-синтаксис.

2. Комбинирование условий с логическими операторами

Условие AND (&):

filtered_df = df.filter((df["age"] > 30) & (df["salary"] > 50000))

Условие OR (|):

filtered_df = df.filter((df["age"] > 30) | (df["department"] == "HR"))

Условие NOT (~):

filtered_df = df.filter(~(df["department"] == "HR"))

3. Фильтрация со строковыми условиями

Методы filter() и where() поддерживают строковые выражения в SQL-стиле:

filtered_df = df.filter("age > 30 AND department = 'HR'")

4. Фильтрация с isin()

Метод isin() фильтрует по набору значений — аналог оператора SQL IN:

# Фильтрация, где отдел 'HR' или 'Finance'
filtered_df = df.filter(df["department"].isin("HR", "Finance"))

5. Фильтрация с startswith()

# Фильтрация строк, где имя начинается с 'A'
filtered_df = df.filter(df["name"].startswith("A"))

6. Фильтрация с endswith()

# Фильтрация строк, где имя заканчивается на 'son'
filtered_df = df.filter(df["name"].endswith("son"))

7. Паттерн-матчинг: like() и rlike()

Базовый паттерн-матчинг с like():

# Имена, начинающиеся с 'A'
filtered_df = df.filter(df["name"].like("A%"))

Regex-паттерн-матчинг с rlike():

# Имена, содержащие 'son'
filtered_df = df.filter(df["name"].rlike("son"))

8. Работа с null: isNull() и isNotNull()

# Фильтрация строк, где age равен NULL
filtered_df = df.filter(df["age"].isNull())

# Фильтрация строк, где age НЕ равен NULL
filtered_df = df.filter(df["age"].isNotNull())

9. Фильтрация через SQL-запросы

# Регистрация DataFrame как временной таблицы
df.createOrReplaceTempView("employees")

# Использование SQL-запроса для фильтрации
filtered_df = spark.sql("SELECT * FROM employees WHERE age > 30")

Полный пример

from pyspark.sql import SparkSession
from pyspark.sql.functions import col

# Инициализация Spark Session
spark = SparkSession.builder \
    .appName("PySpark Filtering Example") \
    .getOrCreate()

# Тестовый DataFrame
data = [
    ("Alice", "HR", 5000),
    ("Bob", "IT", 6000),
    ("Charlie", "Finance", 7000),
    ("David", "IT", 4000),
    ("Eve", "HR", 5500),
    ("Frank", "Finance", 8000),
]
columns = ["name", "department", "salary"]

df = spark.createDataFrame(data, columns)

# Показать исходные данные
print("Original Data:")
df.show()

# Фильтр: зарплата больше 6000
print("Filter: Salary > 6000")
df.filter(df.salary > 6000).show()

# Фильтр: конкретный отдел
print("Filter: Department = 'IT'")
df.filter(df.department == "IT").show()

# Комбинирование условий
print("Filter: Salary > 5000 and Department = 'HR'")
df.filter((df.salary > 5000) & (df.department == "HR")).show()

# Через where()
print("Filter: Salary < 6000 using where()")
df.where("salary < 6000").show()

# Через isin()
print("Filter: Department in ('IT', 'HR')")
df.filter(df.department.isin("IT", "HR")).show()

# Фильтрация по null-значениям
data_with_nulls = [
    ("Alice", None, 5000),
    ("Bob", "IT", 6000),
    (None, "Finance", 7000),
    ("David", "IT", 4000),
]
df_with_nulls = spark.createDataFrame(data_with_nulls, columns)

print("Filter: Rows where department is not null")
df_with_nulls.filter(col("department").isNotNull()).show()

print("Filter: Rows where name is null")
df_with_nulls.filter(col("name").isNull()).show()

# Остановка Spark Session
spark.stop()

Итог

Операции фильтрации можно тонко настраивать с помощью различных методов: от базовых условий до паттерн-матчинга, работы с null и SQL-запросов — всё это позволяет адаптировать фильтры под любые сценарии работы с данными в привычном SQL-стиле.