Фильтрация данных в PySpark
Методы filter() и where(), логические операторы AND/OR/NOT, isin(), startswith(), endswith(), like(), rlike(), isNull(), isNotNull(), SQL-запросы через createOrReplaceTempView.
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-стиле.