DataFrame Wrangling: filtering joins, set operations и wide↔long
Полная таксономия join-семантик, set operations для снапшот-диффов, pivot через DataFrame API и шаблон wide-to-long через posexplode.
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 = NULL → NULL (не 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()- аналог SQLEXCEPT DISTINCT: удаляет все строки из Y, которые встречаются в Z, независимо от количества повторений.exceptAll()- аналог SQLEXCEPT 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:
- Сравнивается вся строка целиком. Если строка с ключом
id=5изменила полеstatus, это будет одновременно DELETE (старая версия) и INSERT (новая версия). Нет понятия UPDATE. subtractиintersectвнутри делают Shuffle и distinct. Это дорого на терабайтных таблицах. Для таких объёмов используйте предварительную фильтрацию по партиции (year,month) или full outer join с ключами.- Дубли внутри снапшота нарушают логику: если одна и та же строка в
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) |