Выборка колонок в PySpark 📋

11 способов выборки колонок: select(), col(), withColumn(), alias(), withColumnRenamed(), expr(), drop(), SQL-синтаксис, динамическая выборка и вложенные структуры.

core

В PySpark выборка колонок из DataFrame — ключевая операция, аналогичная оператору SQL SELECT. Существует множество способов выбирать колонки, что даёт гибкость в обработке и просмотре данных.


1. Базовая выборка с select()

# Выбор одной колонки
selected_df = df.select("name")

# Выбор нескольких колонок
selected_df = df.select("name", "age", "department")

Метод select() создаёт новый DataFrame, содержащий только указанные колонки.


2. Использование col() в выражениях

from pyspark.sql.functions import col

# Выбор колонок через col()
selected_df = df.select(col("name"), col("age"))

Этот подход удобен для операций, требующих явности, особенно когда имена колонок содержат пробелы или специальные символы.


3. Добавление колонок с withColumn()

from pyspark.sql.functions import col, lit

# Добавление новой колонки со статическим значением
df_with_constant = df.withColumn("status", lit("active"))

# Добавление новой колонки на основе существующей
df_with_computed = df.withColumn("age_plus_5", col("age") + 5)

Цепочка нескольких вызовов withColumn():

# Добавление нескольких новых колонок
df_updated = df.withColumn("status", lit("active")) \
               .withColumn("next_year_age", col("age") + 1)

Этот метод особенно полезен, когда нужно обогатить DataFrame вычисленными или дополнительными данными.


4. Переименование колонок с alias()

# Переименование колонок при выборке (аналог SQL AS)
selected_df = df.select(col("name").alias("employee_name"), col("age"))

В результате DataFrame будет содержать колонку employee_name вместо name.


5. Переименование колонок с withColumnRenamed()

# Переименование одной колонки
df_renamed = df.withColumnRenamed("name", "employee_name")

# Переименование нескольких колонок (цепочка)
df_renamed = df.withColumnRenamed("name", "employee_name") \
               .withColumnRenamed("age", "employee_age")

Метод особенно полезен при очистке имён колонок после чтения из «грязных» источников данных.


6. Выборка с выражениями через expr()

from pyspark.sql.functions import expr

# Выборка с выражениями
selected_df = df.select(expr("age + 1 AS next_year_age"), "department")

Этот метод позволяет применять SQL-подобные выражения, что интуитивно понятно для SQL-разработчиков.


7. Выбор всех колонок с исключением ненужных

# Выбор всех колонок, кроме 'salary'
selected_df = df.select("*").drop("salary")

Удобно, когда нужно взять большинство колонок DataFrame, исключив лишь некоторые.


8. Использование SQL-синтаксиса

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

# Выборка колонок через SQL
selected_df = spark.sql("SELECT name, age FROM employees")

Метод удобен для сложных запросов или если вы лучше знакомы с SQL.


9. Динамическая выборка колонок

# Список колонок для выборки
columns_to_select = ["name", "age"]

# Динамическая выборка колонок
selected_df = df.select(*columns_to_select)

Метод даёт гибкость, особенно когда набор колонок может меняться во время выполнения.


10. Исключение колонок с drop()

# Сохранить все колонки, кроме 'salary' и 'address'
selected_df = df.drop("salary", "address")

Удобно при работе с большим числом колонок, когда нужно убрать лишь некоторые из них.


11. Работа с вложенными структурами

# Предположим, что `address` — это struct-колонка с полями `city` и `state`
selected_df = df.select("name", "address.city", "address.state")

Позволяет обращаться к вложенным полям — аналогично доступу к полям JSON-объекта.


Итог

Выборка колонок в PySpark предоставляет множество методов, близких к SQL-синтаксису. От базовых выборок до выражений и динамических методов — у вас есть все инструменты для эффективной работы с DataFrame.