Джоины в PySpark
Типы джоинов (inner, left, right, outer, cross, left_semi, left_anti), синтаксис join(), джоин по нескольким колонкам, условные выражения, выборка колонок после джоина.
core
Обзор¶
Джоины в PySpark аналогичны SQL-джоинам: они позволяют объединять данные из двух и более DataFrame по общей колонке.
Типы джоинов в PySpark¶
Функция join() в PySpark поддерживает следующие типы джоинов:
| Тип джоина | Описание |
|---|---|
inner |
Совпадающие строки из обоих DataFrame |
left |
Все строки из левого DataFrame, null для несовпадений |
right |
Все строки из правого DataFrame, null для несовпадений |
outer |
Все строки из обоих DataFrame, null для несовпадений |
cross |
Декартово произведение обоих DataFrame |
left_semi |
Внутренний джоин, возвращающий только колонки левого DataFrame |
left_anti |
Строки из левого DataFrame, для которых нет совпадений в правом |
Синтаксис и параметры¶
DataFrame.join(other, on=None, how=None)
other: DataFrame, с которым выполняется джоинon: Колонка(и) для джоина (строка, список или выражение)how: Тип джоина ("inner","left","right","outer","cross")
Inner Join¶
result = df1.join(df2, on="id", how="inner")
result.show()
Практические применения: объединение клиентов с заказами, сотрудников с проектами, товаров с данными о продажах.
Left Outer Join¶
result = df1.join(df2, on="id", how="left")
result.show()
Полезен для сохранения всех записей из основного набора данных при добавлении связанной информации там, где она доступна.
Left Anti Join¶
result = df1.join(df2, on="id", how="left_anti")
result.show()
Позволяет найти несвязанные записи: например, клиентов без заказов или товары, которые никогда не продавались.
Джоин по нескольким колонкам¶
result = df1.join(df2, on=["id", "department"], how="inner")
result.show()
Джоин без явного указания on¶
# PySpark автоматически выполняет джоин по колонкам с одинаковыми именами
result = df1.join(df2, how="inner")
result.show()
Использование условий в джоинах¶
from pyspark.sql import functions as F
result = df1.join(df2, df1["id"] == df2["user_id"], how="inner")
result.show()
Джоин с несколькими условиями¶
joined_df = sales_df.join(
customers_df,
(sales_df["customer_id"] == customers_df["customer_id"]) & (sales_df["region"] == customers_df["region"]),
"inner"
)
Выборка конкретных колонок после джоина¶
result_df = sales_df.join(customers_df, on=["customer_id"], how="inner").select(
sales_df["*"],
customers_df["name"].alias("customer_name")
)
Полный пример¶
from pyspark.sql import SparkSession
# Инициализация Spark Session
spark = SparkSession.builder \
.appName("PySpark Join Example") \
.getOrCreate()
# Тестовый DataFrame 1 (сотрудники)
data1 = [
(1, "Alice", "HR"),
(2, "Bob", "IT"),
(3, "Charlie", "Finance"),
(4, "David", "IT")
]
columns1 = ["id", "name", "dept"]
df1 = spark.createDataFrame(data1, columns1)
# Тестовый DataFrame 2 (отделы)
data2 = [
("HR", "Human Resources"),
("IT", "Information Technology"),
("Marketing", "Marketing"),
]
columns2 = ["dept", "dept_name"]
df2 = spark.createDataFrame(data2, columns2)
# Inner Join
inner_join = df1.join(df2, on="dept", how="inner")
print("Inner Join Result:")
inner_join.show()
# Left Join
left_join = df1.join(df2, on="dept", how="left")
print("Left Join Result:")
left_join.show()
# Full Outer Join
outer_join = df1.join(df2, on="dept", how="outer")
print("Full Outer Join Result:")
outer_join.show()
# Остановка Spark Session
spark.stop()
Лучшие практики джоинов в PySpark¶
- Явно указывайте имена колонок через параметр
on - Используйте псевдонимы (
alias) для колонок во избежание конфликтов при совпадающих именах - Применяйте условные выражения, когда имена колонок различаются или требуется дополнительная логика