DataFrame Wrangling: filtering joins, set operations и wide↔long

Полная таксономия join-семантик, set operations для снапшот-диффов, pivot через DataFrame API и шаблон wide-to-long через posexplode.

api wrangling

DataFrame Wrangling: filtering joins, set operations и wide↔long

Data Wrangling - это «борьба с данными»: приведение их из того вида, в котором они приходят из источника, в тот вид, который нужен потребителю. В реальных пайплайнах 70–80% кода приходится именно на трансформации: фильтрацию, обогащение через join, слияние снапшотов, изменение формы таблиц. Правильный выбор семантики каждой операции напрямую влияет на корректность результата и стоимость вычислений.


Две семантики join

Spark-джойны делятся на два класса по тому, что они возвращают:

  • Mutating joins - добавляют колонки из правой таблицы (left, right, inner, full)
  • Filtering joins - только фильтруют строки левой таблицы, не добавляя новых колонок (left_semi, left_anti)

Это разграничение важно архитектурно: когда вам нужно просто проверить наличие записи в другой таблице - не обогащая строку новыми данными - filtering join корректен и эффективен. Mutating join в той же ситуации создаёт лишние колонки, требует select("a.*") и может дублировать строки при неуникальных ключах в правой таблице.


Полная таксономия join-семантик

Разберём все типы джойнов на конкретных примерах. Допустим, таблица A содержит ключи {1, 2, 3}, таблица B - ключи {2, 3, 4}.

Inner Join - самый распространённый. Возвращает только строки, где ключ есть в обеих таблицах. Главная ловушка: если ключи в B не уникальны, строки из A дублируются. Один заказ из A + три строки с тем же product_id в B = три строки в результате. Это не баг Spark - это корректная реляционная семантика, которая часто удивляет новичков.

# Опасно при неуникальных ключах в правой таблице
orders.join(products, on="product_id", how="inner")
# Если у одного product_id в products 3 строки → каждый order дублируется 3 раза

Left Join - сохраняет все строки из левой таблицы. Строки без совпадения получают NULL в колонках из правой таблицы. Стандарт для обогащения данных: даже если справочника нет - запись должна остаться.

Right Join - зеркало left join. Сохраняет все строки из правой таблицы. На практике редко используется напрямую; чаще просто меняют порядок таблиц и делают left join.

Full Outer Join - сохраняет все строки из обеих таблиц. Строки без пары получают NULL с противоположной стороны. Дорого: требует Shuffle по ключу с обеих сторон, результат может быть значительно больше любой из таблиц.

Cross Join - декартово произведение: каждая строка A с каждой строкой B. Результат: |A| × |B| строк. 1000 × 1000 = 1 000 000. Обычно это катастрофа производительности, но иногда нужен: генерация комбинаций дат и магазинов для calendar-плотных рядов, создание всех пар для сравнения. Spark требует явного указания how="cross" (или crossJoin()) как защиту от случайного декартова произведения.

# Явный cross join - Spark не даст случайно его сделать
stores.crossJoin(dates)   # или .join(dates, how="cross")

# Генерация всех комбинаций (store, date) для заполнения пропусков в временных рядах
stores.crossJoin(date_range) \
    .join(actual_sales, on=["store_id", "date"], how="left")

Filtering joins: left_semi и left_anti

left_semi - есть ли совпадение?

Эквивалент SQL WHERE EXISTS (SELECT 1 FROM B WHERE A.key = B.key).

# Оставить только orders, у которых есть запись в active_customers
orders.join(active_customers, on="customer_id", how="left_semi")

Чем отличается от inner join:

  • Не добавляет колонки из правой таблицы - не нужен select(orders["*"])
  • Не дублирует строки левой таблицы, даже если в правой несколько совпадений по ключу
  • Catalyst может оптимизировать агрессивнее: правую таблицу не нужно полностью передавать, достаточно знать, есть ли ключ
# inner join при неуникальных ключах в B даёт дубли:
A.join(B, on="key", how="inner")       # может дать 3 строки из одной строки A
A.join(B, on="key", how="left_semi")   # всегда не более одной строки на строку A

Под капотом. Catalyst транслирует left_semi в LeftSemi физический оператор. При Broadcast Join правая таблица хешируется и рассылается на все ноды. Для каждой строки из левой таблицы достаточно проверить наличие ключа в хеш-таблице - не нужно читать значения правой стороны. Это делает left_semi быстрее эквивалентного inner join + select(a.*) на больших данных.

Практический сценарий. Silver-слой: фильтрация транзакций только по активным пользователям из справочника:

from pyspark.sql import functions as F

transactions = spark.read.parquet("s3://bronze/transactions/")
active_users = spark.read.parquet("s3://silver/users/").filter(F.col("is_active") == True)

# Только транзакции активных пользователей - без добавления колонок из users
clean_transactions = transactions.join(
    active_users.select("user_id"),  # берём только ключ из правой стороны
    on="user_id",
    how="left_semi"
)

left_anti - нет ли совпадения?

Эквивалент SQL WHERE NOT EXISTS. Базовый паттерн для инкрементальной загрузки:

# Загрузить только строки из new_data, которых ещё нет в target
new_data.join(target, on="event_id", how="left_anti") \
    .write.mode("append").save(path)

left_anti надёжнее, чем NOT IN (subquery). Критичная разница: NOT IN с NULL в правой таблице вернёт пустой результат из-за трёхзначной логики SQL:

# SQL семантика: NOT IN с NULL - ловушка
# WHERE id NOT IN (1, 2, NULL)
# = WHERE id != 1 AND id != 2 AND id != NULL
# = WHERE ... AND UNKNOWN  → всегда UNKNOWN → 0 строк!

# left_anti корректно обрабатывает NULL в правой таблице:
new_records.join(existing, on="id", how="left_anti")  # безопасно

Сценарий: Bronze → Silver без дублей. При ежедневной загрузке файлов из S3:

# Приходят данные за сегодня - часть уже была в Silver вчера
daily_raw = spark.read.parquet(f"s3://landing/{today}/")
silver    = spark.read.parquet("s3://silver/events/")

# Только действительно новые события
new_events = daily_raw.join(silver, on="event_id", how="left_anti")

new_events.write.mode("append").parquet("s3://silver/events/")

Поиск «осиротевших» записей (Orphaned Records). Проверка целостности данных: заказы без покупателей, транзакции без счетов:

# Найти транзакции для несуществующих счетов
orphaned = transactions.join(accounts, on="account_id", how="left_anti")

if orphaned.count() > 0:
    orphaned.write.parquet("s3://quarantine/orphaned_transactions/")
    raise DataQualityError(f"Найдено {orphaned.count()} осиротевших транзакций")

NULL-safe joins и edge cases

В Spark SQL, как и в стандарте SQL, NULL = NULLNULL (не True). Это значит: обычный join по ключу, в котором есть NULL, никогда не соединит строки с NULL с обеих сторон - они просто не совпадут.

# Стандартный join: NULL не совпадает с NULL
A = spark.createDataFrame([(None, "a"), (1, "b")], ["id", "val_a"])
B = spark.createDataFrame([(None, "x"), (1, "y")], ["id", "val_b"])

A.join(B, on="id", how="inner").show()
# Результат: только строка (1, "b", "y") - строки с NULL не совпали

Для специфических задач (например, join по колонке, которая намеренно может быть NULL, и нужно соединить NULL с NULL) используйте null-safe equality оператор <=>:

from pyspark.sql import functions as F

# NULL-safe join: NULL совпадает с NULL
A.join(B, A["id"].eqNullSafe(B["id"]), how="inner").show()
# Результат: (None, "a", "x") и (1, "b", "y") - оба совпали

<=> - это Diamond Operator из SQL:2003. В pandas аналог: pd.Series.eq(other, fill_value=...). В Spark это Column.eqNullSafe(other) или SQL a.id <=> b.id.

Когда нужен null-safe join:

  • Слияние двух снапшотов по составному ключу, где одно из полей может быть NULL (например, segment_id не заполнен для новых пользователей)
  • Дедупликация строк с возможными NULL в ключевых полях
  • Сравнение типа "найти строки, где значение изменилось" - NULL → значение и значение → NULL это тоже изменение

Broadcast join в wrangling workloads

При join большой таблицы с маленькой Spark может автоматически разослать маленькую таблицу всем исполнителям - это Broadcast Join. Никакого Shuffle не нужно: каждый исполнитель выполняет join локально по своей партиции большой таблицы.

from pyspark.sql import functions as F

# Автоматический broadcast: если small_table < spark.sql.autoBroadcastJoinThreshold (10MB)
transactions.join(country_codes, on="country_id", how="left")
# Catalyst сам применит BroadcastHashJoin

# Явный hint - надёжнее, когда размер известен
transactions.join(
    F.broadcast(country_codes),
    on="country_id",
    how="left"
)

Пороговое значение для авто-broadcast: spark.sql.autoBroadcastJoinThreshold (дефолт 10 MB). Если справочник больше порога, broadcast не сработает автоматически - нужен явный hint.

Критически важно: F.broadcast() не сортирует и не дедuplicit справочник. Если в маленькой таблице неуникальные ключи - broadcast hash join всё равно даст дубли. Всегда проверяйте уникальность ключей в broadcast-таблице перед применением:

# Безопасный паттерн: дедупликация перед broadcast
clean_ref = reference.dropDuplicates(["key_col"])
assert clean_ref.count() == reference.count(), "Найдены дубли в справочнике!"

transactions.join(F.broadcast(clean_ref), on="key_col", how="left")

Set operations

Для двух DataFrame с идентичной схемой (имена и типы колонок в том же порядке):

# SQL INTERSECT: строки, которые есть в обоих
Y.intersect(Z)

# SQL EXCEPT / MINUS: строки из Y, которых нет в Z
Y.subtract(Z)    # или Y.exceptAll(Z) - см. различие ниже

# SQL UNION ALL: объединение с дублями
Y.union(Z)

# SQL UNION: объединение без дублей
Y.union(Z).dropDuplicates()

intersect и subtract сравнивают всю строку целиком - аналогично dropDuplicates() поверх join. Это делает их удобным инструментом для снапшот-диффов.

subtract (except) vs exceptAll: важная разница

Два метода выглядят похоже, но семантика отличается при дублях:

Y = spark.createDataFrame([(1,), (1,), (2,)], ["id"])
Z = spark.createDataFrame([(1,), (3,)],        ["id"])

Y.subtract(Z).show()   # [2]       - removes ALL occurrences of matching rows
Y.exceptAll(Z).show()  # [1, 2]    - removes ONE occurrence per match
  • subtract() - аналог SQL EXCEPT DISTINCT: удаляет все строки из Y, которые встречаются в Z, независимо от количества повторений.
  • exceptAll() - аналог SQL EXCEPT ALL: удаляет по одному вхождению за каждое совпадение в Z.

В задачах снапшот-диффа обычно нужен subtract - нас не интересуют дубли внутри одного снапшота.

union vs unionByName: безопасная эволюция схем

union() соединяет DataFrame по позиции колонок, игнорируя имена. Это смертельно опасно при эволюции схемы:

# ОПАСНО: union по позиции
schema_v1 = spark.createDataFrame([(1, "Alice", 100)], ["id", "name", "amount"])
schema_v2 = spark.createDataFrame([(2, 200, "Bob")],   ["id", "amount", "name"])

schema_v1.union(schema_v2).show()
# id | name  | amount
# 1  | Alice | 100
# 2  | 200   | Bob      ← name и amount перепутаны! Spark не предупредил.

unionByName() соединяет по именам колонок, что безопасно при разном порядке. Параметр allowMissingColumns=True (PySpark 3.1+) автоматически заполняет NULL для отсутствующих колонок:

# БЕЗОПАСНО: unionByName по именам
schema_v1.unionByName(schema_v2).show()
# Правильный результат - колонки выровнены по именам

# allowMissingColumns=True для эволюции схемы
old_data = spark.read.parquet("s3://bucket/2024/")  # без колонки "region"
new_data = spark.read.parquet("s3://bucket/2025/")  # с колонкой "region"

combined = old_data.unionByName(new_data, allowMissingColumns=True)
# Строки из old_data получат NULL в колонке "region"

Правило: в production-пайплайнах всегда используйте unionByName. union() оправдан только если схема заведомо одинакова и это явно задокументировано.


Паттерн Snapshot Diff: расчёт инкремента через set operations

Один из ключевых паттернов Data Engineering - вычисление изменений между двумя снапшотами данных. Это нужно для:

  • SCD Type 2 (медленно меняющиеся измерения) - определить, что изменилось
  • Инкрементальной загрузки - загрузить только новые и изменённые записи
  • Аудита - показать, что появилось или исчезло между двумя периодами

Полный пример снапшот-диффа с тегами типа изменения:

from pyspark.sql import functions as F

today     = spark.read.parquet("s3://bucket/snapshot/2026-05-16/")
yesterday = spark.read.parquet("s3://bucket/snapshot/2026-05-15/")

# Что появилось - есть сегодня, не было вчера
added   = today.subtract(yesterday).withColumn("change_type", F.lit("INSERT"))

# Что исчезло - было вчера, нет сегодня
removed = yesterday.subtract(today).withColumn("change_type", F.lit("DELETE"))

# Что не изменилось
unchanged = today.intersect(yesterday).withColumn("change_type", F.lit("UNCHANGED"))

# Инкремент для записи в Delta Lake
delta_log = added.unionByName(removed)

delta_log.write \
    .mode("append") \
    .partitionBy("change_type") \
    .parquet("s3://silver/change_log/")

Ограничения подхода через set operations:

  1. Сравнивается вся строка целиком. Если строка с ключом id=5 изменила поле status, это будет одновременно DELETE (старая версия) и INSERT (новая версия). Нет понятия UPDATE.
  2. subtract и intersect внутри делают Shuffle и distinct. Это дорого на терабайтных таблицах. Для таких объёмов используйте предварительную фильтрацию по партиции (year, month) или full outer join с ключами.
  3. Дубли внутри снапшота нарушают логику: если одна и та же строка в today встречается дважды, intersect вернёт её один раз.

Альтернатива через full outer join для управления UPDATE-случаями:

# Более гибкий подход - видим конкретный ключ и обе версии
diff = today.alias("new").join(
    yesterday.alias("old"),
    on="record_id",
    how="full"
)

# Категоризация изменений по статусу совпадения
diff.withColumn(
    "change_type",
    F.when(F.col("old.record_id").isNull(), "INSERT")
     .when(F.col("new.record_id").isNull(), "DELETE")
     .when(
         # Любое поле изменилось - это UPDATE
         F.col("new.status") != F.col("old.status"),
         "UPDATE"
     )
     .otherwise("UNCHANGED")
)

DataFrame Pivot

groupBy().pivot().agg() транспонирует значения одной колонки в заголовки новых колонок. Это перевод из Long-формата в Wide, который часто нужен для BI-витрин и отчётов.

# Исходные данные: store_id, month, revenue
df.groupBy("store_id").pivot("month").sum("revenue").show()
# → store_id | 2026-01 | 2026-02 | 2026-03 | ...

Когда список значений известен заранее - передайте его явно: Spark иначе делает дополнительный pass по данным чтобы собрать distinct значения.

months = ["2026-01", "2026-02", "2026-03"]
df.groupBy("store_id").pivot("month", months).sum("revenue")

Pivot с несколькими агрегатами - колонки получают суффикс:

from pyspark.sql import functions as F

df.groupBy("region").pivot("category", ["electronics", "clothing"]).agg(
    F.sum("qty").alias("qty"),
    F.avg("price").alias("avg_price")
)
# → region | electronics_qty | electronics_avg_price | clothing_qty | clothing_avg_price

Обратная операция (long → wide уже сделана, нужно wide → long) - следующий раздел.

Производительность pivot: динамический vs статический

Это самая распространённая ловушка при работе с pivot. Понять её важно до того, как пайплайн уйдёт в production.

Что происходит при динамическом pivot. Spark не знает заранее, какие значения появятся в колонке month. Чтобы определить имена будущих колонок (схему результата), Spark вынужден сначала выполнить отдельный Job - полный scan данных для сбора DISTINCT month. После этого запускается второй Job - собственно агрегация. Итог: два Job вместо одного, двойное время выполнения, двойной I/O.

В Spark UI это выглядит как два отдельных Job: первый маленький (только DISTINCT), второй - основная агрегация. Новичок смотрит и не понимает, откуда взялся лишний Job.

Решение: всегда передавайте список значений явно. Если значения предсказуемы (месяцы, регионы, статусы) - получите их один раз и используйте константой:

# Получить список значений один раз из метаданных/конфига
REGIONS = ["North", "South", "East", "West"]
MONTHS  = [f"2026-{m:02d}" for m in range(1, 13)]

# Статический pivot - только один Spark Job
sales.groupBy("store_id") \
    .pivot("region", REGIONS) \
    .agg(F.sum("revenue").alias("revenue"))

# Если список значений нужно получить из данных - кешируйте его
all_months = df.select("month").distinct().rdd.flatMap(lambda x: x).collect()
df.groupBy("store_id").pivot("month", all_months).sum("revenue")

Ограничения pivot. Широкий pivot (много уникальных значений) создаёт DataFrame с сотнями колонок. Это проблема:

  • Spark хранит метаданные схемы в Driver - 1000 колонок = 1000 объектов в памяти Driver
  • BI-инструменты плохо работают с таблицами-«простынями»
  • JOIN такой таблицы с другой = O(cols_a × cols_b) по памяти при Broadcast

Если уникальных значений > 50–100 - задумайтесь, нужен ли pivot вообще, или лучше оставить Long-формат.


Wide → Long: posexplode и struct-explode

Задача: широкая таблица с несколькими однотипными метриками (колонки-кварталы, колонки-регионы) нужно привести к нормализованному Long-формату.

posexplode - explode с индексом позиции

explode() разворачивает массив в строки. posexplode() добавляет колонку с порядковым номером элемента, что позволяет восстановить исходный порядок:

from pyspark.sql import functions as F

df = spark.createDataFrame([
    ("user_1", [10.0, 20.0, 30.0]),
    ("user_2", [40.0, 50.0, 60.0]),
], ["user_id", "scores"])

df.select("user_id", F.posexplode("scores").alias("idx", "score")).show()
# user_id | idx | score
# user_1  |  0  | 10.0
# user_1  |  1  | 20.0
# user_1  |  2  | 30.0

posexplode нужен когда:

  • Два параллельных массива нужно соотнести по позиции (feature_names[i] + feature_values[i])
  • Нужно восстановить порядок после explode (порядок строк в Spark не гарантирован без orderBy)
# Соотнесение двух массивов по позиции - без posexplode это невозможно
df.select(
    "user_id",
    F.posexplode("feature_values").alias("pos", "value")
).join(
    df.select("user_id", F.posexplode("feature_names").alias("pos", "name")),
    on=["user_id", "pos"]
)

Почему не просто zip_with? В Spark 2.x zip_with отсутствовал. posexplode работает начиная с ранних версий и хорошо известен командам. В Spark 3.x появился arrays_zip(), который удобнее для попарного соединения массивов:

# Современный подход через arrays_zip (Spark 3.x)
df.withColumn(
    "zipped",
    F.arrays_zip("feature_names", "feature_values")
).select(
    "user_id",
    F.explode("zipped").alias("pair")
).select(
    "user_id",
    F.col("pair.feature_names").alias("name"),
    F.col("pair.feature_values").alias("value")
)

to_long - превращение широкой таблицы в long-формат

Классический паттерн для wide→long, когда набор колонок не известен заранее или слишком велик для ручного select:

from pyspark.sql.functions import explode, array, struct, lit, col

def to_long(df, by: list):
    """Транспонирует все колонки кроме `by` в пары metric/value."""
    value_cols = [c for c in df.columns if c not in by]

    # Все value-колонки должны быть одного типа (ограничение Spark)
    types = {t for c, t in df.dtypes if c in value_cols}
    assert len(types) == 1, f"Разные типы в value-колонках: {types}"

    kvs = explode(array([
        struct(lit(c).alias("metric"), col(c).alias("value"))
        for c in value_cols
    ])).alias("kvs")

    return df.select(by + [kvs]).select(by + ["kvs.metric", "kvs.value"])

Пример использования:

wide = spark.createDataFrame([
    ("2026-01-01", 1500.0, 320.0, 88.5),
    ("2026-01-02", 1620.0, 305.0, 91.2),
], ["date", "revenue", "costs", "margin_pct"])

to_long(wide, by=["date"]).show()
# date       | metric     | value
# 2026-01-01 | revenue    | 1500.0
# 2026-01-01 | costs      |  320.0
# 2026-01-01 | margin_pct |   88.5

Почему array(struct(...)) - правильный подход:

  • Это распределённая операция - каждый исполнитель трансформирует свои строки без Shuffle
  • Catalyst понимает структуру: может оптимизировать Project → Explode
  • В отличие от stack() (SQL-функция) - строго типизированно, Spark знает схему до выполнения

Почему не SQL stack()? stack(N, expr1, expr2, ...) требует хардкода числа метрик. При добавлении новой колонки нужно менять SQL-строку. to_long() через Python-генератор автоматически подхватывает все value-колонки.

Зачем нужен long-формат:

  • Передача в агрегации по метрикам: groupBy("metric").agg(F.avg("value"))
  • Снапшот-диффы через subtract: сравниваем строки (date, metric, value) целиком
  • Запись в ClickHouse/InfluxDB, которые ожидают long-формат для временных рядов
  • Универсальные дашборды: одна таблица с (date, metric_name, value) вместо десятков колонок

Дедупликация и distinct-семантики

Дедупликация - частая задача в Data Engineering: источники часто присылают дубли из-за at-least-once гарантий доставки, ретраев, или просто ошибок в ETL.

distinct() vs dropDuplicates():

# distinct: убрать полностью одинаковые строки (по ВСЕМ колонкам)
df.distinct()

# dropDuplicates: убрать дубли по подмножеству колонок (сохраняет первую встреченную)
df.dropDuplicates(["event_id"])             # по одной колонке
df.dropDuplicates(["user_id", "date"])      # по составному ключу

Важный нюанс: dropDuplicates(subset) в Spark не детерминирован - при нескольких строках с одинаковым ключом будет сохранена случайная (зависит от порядка обработки партиций). Если нужно сохранить "последнюю" или "первую" версию - используйте Window:

from pyspark.sql.window import Window
from pyspark.sql import functions as F

# Детерминированная дедупликация: оставить последнюю по времени
w = Window.partitionBy("event_id").orderBy(F.col("created_at").desc())

deduped = df.withColumn("rn", F.row_number().over(w)) \
            .filter(F.col("rn") == 1) \
            .drop("rn")

Производительность. И distinct(), и dropDuplicates() внутри делают HashAggregate + Shuffle. Это дорого. Если дубли возникают только в пределах одного файла/партиции (например, дважды отправили один и тот же файл) - можно сначала дедублировать локально:

# Дедупликация только в пределах каждой партиции - без Shuffle
# Осторожно: только если дубли точно не пересекают границы партиций
df.mapInPandas(
    lambda dfs: (pdf.drop_duplicates(subset=["event_id"]) for pdf in dfs),
    schema=df.schema
)

Стратифицированная выборка: sampleBy

sample(fraction) берёт случайную долю от всего DataFrame. Если нужно сохранить распределение по категориям или скорректировать дисбаланс классов - используйте sampleBy:

# Стратифицированная выборка с коррекцией дисбаланса
fractions = {
    1: 0.10,   # 10% от класса "positive" (их мало - берём почти всё)
    0: 0.02,   # 2%  от класса "negative" (их много - прореживаем)
}
df.sampleBy(col="label", fractions=fractions, seed=42)

sampleBy работает без замены (without replacement). Полезен для:

  • Создания dev-выборки перед randomSplit, сохраняя распределение
  • Быстрого профилирования на 1–5% данных без риска потерять редкие значения
  • Balancing классов перед ML без SMOTE и oversampling
# Сначала sampleBy чтобы сбалансировать, затем randomSplit для train/test
balanced = df.sampleBy("label", {0: 0.01, 1: 0.5}, seed=0)
train, test = balanced.randomSplit([0.8, 0.2], seed=42)

Антипаттерны DataFrame Wrangling

Ошибки, которые допускают даже опытные разработчики:

1. union() вместо unionByName() при эволюции схемы.

При добавлении новой колонки в один из источников union() тихо перепутает значения. Пайплайн не упадёт, но данные будут неправильными. Всегда unionByName().

2. Динамический pivot без явного списка значений.

Скрытый дополнительный Job = двойное время выполнения. Всегда передавайте список значений явно или соберите его один раз через collect() и зафиксируйте.

3. Inner join при необходимости left join.

Если справочник неполный (не все ключи из левой таблицы присутствуют в правой), inner join молча удалит строки. Вместо ошибки - тихая потеря данных. Всегда думайте о "ожидаемом" количестве строк на выходе.

4. NOT IN вместо left_anti.

NOT IN с NULL в правой таблице возвращает пустой результат - неочевидная ловушка SQL. Всегда используйте left_anti.

5. join без проверки уникальности ключей.

# Неявный взрыв кардинальности при неуникальных ключах
orders.join(payments, on="order_id", how="inner")
# Если в payments несколько платежей на один заказ - дубли!

# Безопасно: дедупликация перед join
payments_deduped = payments.dropDuplicates(["order_id"])
orders.join(payments_deduped, on="order_id", how="inner")

6. subtract/intersect на терабайтных таблицах без фильтрации.

subtract делает Shuffle + distinct по всем данным. На терабайтах - многочасовой Job. Всегда фильтруйте по партиции (date, region) до операции.

7. Nested explode без ограничения глубины.

# Взрыв массива внутри массива - O(n²) строк
df.select(F.explode("outer_array")).select(F.explode("inner_array"))
# Добавьте filter() между explode() чтобы срезать ненужные строки

Wrangling в medallion architecture

Разные типы трансформаций характерны для разных слоёв медальон-архитектуры:

Bronze (Raw) - минимальные трансформации. union, unionByName для склейки файлов из разных источников, добавление метаданных (_ingestion_ts, _source_file). Никакой бизнес-логики, никаких join.

Bronze → Silver (Cleansing) - left_anti для дедупликации по ключу между свежей порцией и уже загруженными данными. left_semi для фильтрации строк по справочнику активных сущностей. Исправление типов, заполнение NULL, валидация.

Silver → Gold (Business Logic) - join для обогащения (заказы + покупатели + товары), pivot для построения витрин в Wide-формате, to_long для аналитических агрегатов. Снапшот-диффы через subtract/intersect для SCD Type 2.

Gold → Serving - финальный pivot для BI-инструментов, groupBy().pivot().agg() для отчётных таблиц, sampleBy для экспорта тестовых выборок в аналитические инструменты.


Когда что использовать

Задача Инструмент
Строки A с совпадением в B (без колонок B) left_semi
Инкрементальная загрузка - только новые left_anti
Поиск осиротевших записей left_anti
Снапшот-диффы - что появилось / исчезло subtract / intersect
Объединение двух выборок без дублей union().dropDuplicates()
Объединение с разным порядком колонок unionByName()
Объединение при эволюции схемы unionByName(allowMissingColumns=True)
Строки в колонки с агрегацией groupBy().pivot().agg()
Pivot с известными значениями pivot("col", ["a", "b", "c"])
Массив → строки с индексом позиции posexplode()
Несколько колонок → строки key/value to_long() через explode(array([struct(...)]))
Два массива объединить по позиции arrays_zip() + explode()
Выборка с сохранением распределения классов sampleBy()
Детерминированная дедупликация row_number().over(w).filter(rn == 1)
JOIN с NULL в ключах eqNullSafe() / <=>
Ускорить join большое × маленькое F.broadcast(small_df)