Джоины в 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) для колонок во избежание конфликтов при совпадающих именах
  • Применяйте условные выражения, когда имена колонок различаются или требуется дополнительная логика