Conditional Columns: when/otherwise, coalesce, nullif и NULL-безопасные сравнения
Условная логика в DataFrame API: CASE WHEN, coalesce, nullif, NaN vs NULL, NULL-safe equality и трёхзначная логика SQL в production ETL.
Conditional Columns: when/otherwise, coalesce, nullif и NULL-безопасные сравнения¶
Условная логика - основа любого ETL pipeline. Классификация клиентов по сегментам, выбор лучшего источника для поля email, замена sentinel-значений на настоящий NULL, расчёт флагов активности - всё это реализуется через conditional expressions в DataFrame API.
При этом главный источник ошибок в условной логике - NULL. В SQL NULL - это не «пустое значение», а «неизвестно». Это третье логическое состояние, которое ломает наивные сравнения, фильтры и join-ы. В этом уроке мы разберём как NULL работает изнутри и как правильно с ним обращаться в production.
NULL semantics: трёхзначная логика SQL¶
Прежде чем разбирать when и coalesce, нужно понять фундамент - как SQL обращается с NULL.
В математике логики есть два значения: TRUE и FALSE. В SQL - три: TRUE, FALSE и NULL (UNKNOWN). Это называется трёхзначной логикой (three-valued logic, 3VL):
Самые важные следствия для data engineers:
NULL == NULLвозвращаетNULL, а неTRUE- поэтому join по NULL-ключам теряет строкиFALSE AND NULL=FALSE- это исключение из «NULL заражает всё», потому что результат известен независимо от NULLTRUE OR NULL=TRUE- по той же причине- В
filter(condition)строки проходят только если условие ==TRUE; строки сNULLв условии отбрасываются
df = spark.createDataFrame([(1, None), (2, 100), (3, -5)], ["id", "amount"])
# NULL не проходит ни через > ни через <=
df.filter(col("amount") > 0).show() # только id=2
df.filter(col("amount") <= 0).show() # только id=3
# id=1 (amount=NULL) пропадает в обоих случаях!
# Чтобы включить NULL - добавляй явную проверку
df.filter(col("amount").isNull() | (col("amount") > 0)).show() # id=1 и id=2
when / otherwise: CASE WHEN в DataFrame API¶
when(condition, value) - это Spark-эквивалент SQL CASE WHEN. Условия проверяются сверху вниз, первое совпавшее побеждает. .otherwise(default) задаёт значение для строк, не попавших ни в одну ветку:
from pyspark.sql.functions import when, col
df = spark.createDataFrame([
(1, 1500.0, "COMPLETED"),
(2, 300.0, "PENDING"),
(3, 50.0, "CANCELLED"),
(4, None, "COMPLETED"),
(5, 2000.0, None),
], ["order_id", "amount", "status"])
df.withColumn(
"segment",
when(col("amount") > 1000, "enterprise")
.when(col("amount") > 500, "large")
.when(col("amount") > 100, "medium")
.when(col("amount").isNotNull(), "small")
.otherwise("unknown") # amount IS NULL или не попало ни в одну ветку
).show()
# +--------+------+---------+---------+
# |order_id|amount| status| segment|
# +--------+------+---------+---------+
# | 1|1500.0|COMPLETED|enterprise|
# | 2| 300.0| PENDING| medium|
# | 3| 50.0|CANCELLED| small|
# | 4| null|COMPLETED| unknown| # NULL не проходит через > сравнения
# | 5|2000.0| null|enterprise|
# +--------+------+---------+---------+
Типизация ветвей¶
Все ветки when и otherwise должны возвращать совместимые типы. Spark пытается автоматически привести типы, но лучше контролировать это явно:
# Автоматическое приведение - может дать неожиданные результаты
when(col("flag"), 1).otherwise(0.0) # IntegerType и DoubleType → DoubleType (OK)
when(col("flag"), "active").otherwise(0) # StringType и IntegerType → StringType ("0")
# Явное приведение - надёжно
from pyspark.sql.functions import lit
when(col("flag"), lit(1).cast("double")).otherwise(lit(0).cast("double"))
Порядок условий и производительность¶
Catalyst оптимизирует when-цепочки, но порядок условий всё равно имеет значение для читаемости и корректности. Ставьте наиболее частые или наиболее ограничивающие условия первыми - это аналог short-circuit evaluation:
# Неэффективно: сначала дорогое regexp, потом дешёвое isNull
when(col("email").rlike(r"^[\w.+-]+@[\w-]+\.[a-z]{2,}$"), "valid_email")
.when(col("email").isNull(), "missing")
.otherwise("invalid")
# Лучше: дешёвые проверки (isNull) - вперёд
when(col("email").isNull(), "missing")
.when(col("email").rlike(r"^[\w.+-]+@[\w-]+\.[a-z]{2,}$"), "valid_email")
.otherwise("invalid")
Conditional aggregation: when внутри agg¶
when можно использовать внутри агрегатных функций для условного подсчёта:
from pyspark.sql.functions import sum as fsum, count, when, col
# Разбивка выручки по статусу внутри одного groupBy
df.groupBy("region").agg(
fsum("amount").alias("total_revenue"),
fsum(when(col("status") == "COMPLETED", col("amount")).otherwise(0)).alias("completed_revenue"),
count(when(col("status") == "CANCELLED", 1)).alias("cancelled_count"),
# count(when(...)) считает только не-NULL значения
).show()
count(when(condition, 1)) считает строки где условие истинно - потому что when без otherwise возвращает NULL при ложном условии, а count не считает NULL.
coalesce: выбор первого не-NULL значения¶
coalesce(col1, col2, ..., colN) возвращает первое не-NULL значение из списка аргументов. Это стандартный инструмент для fallback-логики при работе с данными из нескольких источников:
from pyspark.sql.functions import coalesce, col, lit
customers = spark.createDataFrame([
(1, None, "ivan@app.ru", None),
(2, "anna@crm.ru", None, "anna@legacy.ru"),
(3, None, None, None),
(4, "bob@crm.ru", "bob@app.ru", "bob@legacy.ru"),
], ["id", "email_crm", "email_app", "email_legacy"])
customers.select(
col("id"),
# Выбрать лучший email: CRM приоритетнее app, legacy - последний резерв
coalesce(
col("email_crm"),
col("email_app"),
col("email_legacy"),
lit("no-email@unknown.com") # финальный дефолт - не NULL
).alias("final_email"),
).show()
# +---+----------------------+
# | id| final_email|
# +---+----------------------+
# | 1| ivan@app.ru | # CRM null, взяли app
# | 2| anna@crm.ru | # CRM есть, взяли его
# | 3|no-email@unknown.com | # все null, взяли дефолт
# | 4| bob@crm.ru | # CRM есть, взяли его (app проигнорирован)
# +---+----------------------+
coalesce как замена when для простых случаев¶
Частый паттерн - заменить NULL на дефолтное значение. coalesce читается чище чем эквивалентный when:
# Эквивалентные выражения:
when(col("amount").isNull(), lit(0.0)).otherwise(col("amount"))
coalesce(col("amount"), lit(0.0)) # чище и быстрее компилируется
# Заполнение NULL в числовых колонках
df.select(
coalesce(col("price"), lit(0.0)).alias("price"),
coalesce(col("discount"), lit(0.0)).alias("discount"),
coalesce(col("qty"), lit(1)).alias("qty"),
)
coalesce в join по nullable ключу¶
Распространённый сценарий: join по полю, которое может быть NULL в одном из источников, и нужно использовать альтернативный ключ:
# Найти заказы, используя order_id; если нет - по transaction_id
orders.join(
payments,
coalesce(orders["order_id"], orders["transaction_id"]) ==
coalesce(payments["order_id"], payments["transaction_id"]),
"left"
)
nullif: превращение «мусора» в NULL¶
nullif(col, value) возвращает NULL если колонка равна value, иначе - саму колонку. Это стандартный способ превратить sentinel-значения (пустые строки, -1, "UNKNOWN", "N/A") в честный NULL:
from pyspark.sql.functions import nullif, col, lit
df = spark.createDataFrame([
(1, "ivan@mail.ru", "79001234567", "ACTIVE"),
(2, "", "-1", "UNKNOWN"),
(3, "N/A", "79007654321", ""),
(4, None, None, "INACTIVE"),
], ["id", "email", "phone", "status"])
df.select(
col("id"),
nullif(col("email"), lit("")).alias("email"), # "" → NULL
nullif(col("email"), lit("N/A")).alias("email2"),# "N/A" → NULL
nullif(col("phone"), lit("-1")).alias("phone"), # "-1" → NULL
nullif(col("status"), lit("UNKNOWN")).alias("status_clean"),
nullif(col("status"), lit("")).alias("status_ne"),
).show()
# +---+------------+-----------+--------+------------+-----------+
# | id| email| email2| phone|status_clean| status_ne |
# +---+------------+-----------+--------+------------+-----------+
# | 1|ivan@mail.ru|ivan@mail.ru|7900...| ACTIVE| ACTIVE|
# | 2| null| | null| null| UNKNOWN|
# | 3| N/A| null| 7900...| UNKNOWN| null|
# | 4| null| null| null| INACTIVE| INACTIVE|
# +---+------------+-----------+--------+------------+-----------+
Зачем превращать в NULL? Агрегатные функции (avg, sum, count) игнорируют NULL. Пустая строка "" - это не NULL, она будет учтена в count(*) и испортит avg. После nullif данные ведут себя корректно:
# Без nullif: avg включает пустые строки как "числа" (для числовых) или ломается
df.agg({"phone": "count"}).show() # считает все, включая "-1"
# С nullif: только настоящие значения
df.select(nullif(col("phone"), lit("-1")).alias("phone")) \
.agg({"phone": "count"}).show() # только реальные номера
Цепочка nullif через coalesce¶
Для нескольких sentinel-значений объединяем nullif в coalesce:
from pyspark.sql.functions import nullif, coalesce, col, lit
# Превратить в NULL если: пустая строка, "N/A", "UNKNOWN", "NULL" (строка)
def clean_string(c):
return coalesce(
nullif(nullif(nullif(nullif(c, lit("")), lit("N/A")), lit("UNKNOWN")), lit("NULL")),
lit(None)
)
df.select(clean_string(col("email")).alias("email_clean"))
NaN vs NULL: разные понятия для числовых данных¶
NaN (Not a Number) - специальное значение типа double по стандарту IEEE 754. Это не то же самое, что NULL:
from pyspark.sql.functions import isnan, isnull, nanvl, col
import math
df = spark.createDataFrame([
(1, 1.0),
(2, float('nan')),
(3, None),
(4, float('inf')),
], ["id", "value"])
df.select(
col("id"),
col("value"),
isnan(col("value")).alias("is_nan"), # True только для NaN
isnull(col("value")).alias("is_null"), # True только для None/NULL
).show()
# +---+--------+------+-------+
# | id| value|is_nan|is_null|
# +---+--------+------+-------+
# | 1| 1.0| false| false|
# | 2| NaN| true| false| # NaN ≠ NULL!
# | 3| null| false| true|
# | 4|Infinity| false| false|
# +---+--------+------+-------+
NaN ведёт себя иначе чем NULL в агрегациях:
# avg с NaN - результат NaN (заражает всё!)
df.agg({"value": "avg"}).show() # → NaN
# avg с NULL - NULL игнорируется, результат корректный
# Чтобы исправить: заменить NaN на NULL через nanvl
from pyspark.sql.functions import nanvl, lit
df.select(nanvl(col("value"), lit(None)).alias("value_clean")) \
.agg({"value_clean": "avg"}).show() # корректный среднее без NaN
nanvl(col, replacement) возвращает replacement если значение NaN, иначе - само значение. Это аналог coalesce но для NaN.
Практическое правило: при чтении данных из CSV/JSON с числовыми полями всегда проверяйте наличие NaN:
# При ingestion: заменяем NaN на NULL сразу
from pyspark.sql.functions import nanvl, lit
numeric_cols = ["price", "discount", "qty", "score"]
for c in numeric_cols:
df = df.withColumn(c, nanvl(col(c), lit(None)))
NULL-безопасные сравнения¶
Проблема обычного равенства¶
Как мы видели, NULL == NULL возвращает NULL, а не TRUE. Это критично для join-ов и дедупликации:
df = spark.createDataFrame([
(1, "A", "A"),
(2, "B", None),
(3, None, None),
(4, None, "C"),
], ["id", "col_a", "col_b"])
# Обычное сравнение: строка 3 (NULL==NULL) даёт NULL, не True
df.filter(col("col_a") == col("col_b")).show()
# +---+-----+-----+
# | id|col_a|col_b|
# +---+-----+-----+
# | 1| A| A| # только эта строка!
# +---+-----+-----+
# Строки 2, 3, 4 потеряны - NULL в условии = строка выбрасывается
eqNullSafe: NULL-safe equality¶
col_a.eqNullSafe(col_b) (или оператор <=> в Spark SQL) работает как обычное равенство, но NULL <=> NULL возвращает TRUE:
# NULL-safe сравнение: NULL == NULL → True
df.filter(col("col_a").eqNullSafe(col("col_b"))).show()
# +---+-----+-----+
# | id|col_a|col_b|
# +---+-----+-----+
# | 1| A| A| # A == A
# | 3| null| null| # NULL <=> NULL → True!
# +---+-----+-----+
# Строки 2 и 4: B <=> null → False, null <=> C → False
В SQL-синтаксисе через expr:
df.filter(expr("col_a <=> col_b")) # то же самое через SQL-оператор <=>
Применение в Change Data Capture¶
eqNullSafe незаменим в CDC-пайплайнах для определения изменений: нужно сравнить текущее и предыдущее значение, при этом оба могут быть NULL:
# Найти строки, где значение изменилось (включая переходы NULL → значение и значение → NULL)
changes = current.join(previous, "id", "left") \
.filter(~col("current.status").eqNullSafe(col("previous.status")))
# Строки где оба NULL пройдут eqNullSafe → True → ~ → False → не изменились (корректно)
# Строки NULL → "ACTIVE" пройдут eqNullSafe → False → ~ → True → изменились (корректно)
isNull и isNotNull: правильные проверки на NULL¶
Никогда не используйте col == None для проверки на NULL - это не работает (результат NULL, а не True). Правильный способ:
# НЕПРАВИЛЬНО: col == None → всегда NULL
df.filter(col("amount") == None) # ← не работает!
df.filter(col("amount") != None) # ← не работает!
# ПРАВИЛЬНО
df.filter(col("amount").isNull())
df.filter(col("amount").isNotNull())
# В SQL-строках
df.filter("amount IS NULL")
df.filter("amount IS NOT NULL")
NULL propagation: как NULL распространяется¶
Большинство встроенных функций Spark следует правилу: если любой аргумент NULL - результат NULL. Это называется NULL propagation:
Агрегатные функции - важное исключение из правила: sum, avg, min, max игнорируют NULL и считают только реальные значения. count(*) считает все строки, count(col) - только не-NULL значения в колонке.
Булева логика: скобки обязательны¶
Операторы & (AND), | (OR), ~ (NOT) в PySpark имеют нестандартный приоритет по сравнению с операторами сравнения. Без скобок код компилируется, но делает не то что ожидается:
# ОШИБКА: Python вычисляет как col("amount") > (100 | col("status")),
# потому что | имеет более высокий приоритет чем >
df.filter(col("amount") > 100 | col("status") == "COMPLETED")
# ПРАВИЛЬНО: каждое условие в скобках
df.filter(
(col("amount") > 100) | (col("status") == "COMPLETED")
)
# Сложная логика - выносим в именованные переменные
is_high_value = col("amount") > 500
is_completed = col("status") == "COMPLETED"
is_new_customer = col("days_since_signup") < 30
has_email = col("email").isNotNull()
# Теперь условие читается как бизнес-правило
priority_filter = (is_high_value & is_completed) | (is_new_customer & has_email)
df.filter(priority_filter)
Catalyst optimization для условных выражений¶
Catalyst оптимизирует when-выражения несколькими способами:
# Constant folding: константные условия вычисляются один раз
when(lit(True), col("amount")).otherwise(0)
# Catalyst: заменяет на просто col("amount")
# Boolean simplification
when(col("a").isNotNull() & col("a").isNotNull(), col("a"))
# Catalyst: убирает дублирующееся условие
# CombineFilters: несколько filter объединяются
df.filter(col("a") > 0).filter(col("b") > 0)
# Physical plan: одна Filter node с (a > 0 AND b > 0)
when-цепочки в физическом плане компилируются в JVM-байткод через WholeStageCodeGen - всё в одном проходе по строке без промежуточных объектов. Python if/elif/else в UDF - отдельный процесс с pickle-сериализацией.
Bronze → Silver: нормализация профилей клиентов¶
Полный пример нормализации сырых данных с conditional logic:
from pyspark.sql.functions import (
col, when, coalesce, nullif, nanvl, isnan,
lit, trim, lower, regexp_extract, length
)
raw_customers = spark.read.json("/data/bronze/customers/")
silver_customers = (
raw_customers
# --- Очистка email: три источника, sentinel-значения ---
.withColumn(
"email_clean",
coalesce(
# nullif убирает пустые строки и "N/A" перед coalesce
nullif(trim(lower(col("email_crm"))), lit("")),
nullif(trim(lower(col("email_app"))), lit("N/A")),
nullif(trim(lower(col("email_legacy"))), lit("")),
)
)
# --- Телефон: убрать sentinel -1 ---
.withColumn(
"phone_clean",
nullif(col("phone"), lit("-1"))
)
# --- NaN в числовых полях → NULL ---
.withColumn("score", nanvl(col("score"), lit(None)))
.withColumn("amount", nanvl(col("amount"), lit(None)))
# --- Сегмент клиента ---
.withColumn(
"segment",
when(col("total_spent") > 100_000, "vip")
.when(col("total_spent") > 10_000, "premium")
.when(col("total_spent") > 1_000, "regular")
.when(col("total_spent").isNotNull(), "new")
.otherwise("unknown")
)
# --- Флаг активности: сложная логика ---
.withColumn(
"is_active",
(col("days_since_login") <= 30) &
col("email_clean").isNotNull() &
(col("status") != "BLOCKED")
)
# --- Качество записи ---
.withColumn(
"data_quality",
when(
col("email_clean").isNotNull() &
col("phone_clean").isNotNull() &
col("first_name").isNotNull(),
"complete"
).when(
col("email_clean").isNotNull() | col("phone_clean").isNotNull(),
"partial"
).otherwise("poor")
)
)
Best practices: когда условий становится много¶
Если when-цепочка разрастается до 20+ условий (например, маппинг кода страны → регион), лучше использовать справочник через Broadcast Join:
# АНТИПАТТЕРН: гигантский when для маппинга значений
when(col("country_code") == "RU", "Russia")
.when(col("country_code") == "BY", "Belarus")
.when(col("country_code") == "KZ", "Kazakhstan")
# ... ещё 50 строк ...
.otherwise("Other")
# ЛУЧШЕ: справочная таблица + broadcast join
country_dict = spark.createDataFrame([
("RU", "Russia"), ("BY", "Belarus"), ("KZ", "Kazakhstan"),
], ["code", "name"])
df.join(broadcast(country_dict), df["country_code"] == country_dict["code"], "left") \
.withColumn("country_name", coalesce(col("name"), lit("Other"))) \
.drop("code", "name")
Broadcast join с маленькой таблицей (< 10 МБ) не требует shuffle - справочник рассылается каждому executor и join выполняется локально. Это быстрее 50-уровневого when и проще обновлять.
Практика¶
Задание 1: Нормализация профиля клиента¶
Дан DataFrame с фрагментированными данными о клиентах:
customers = spark.createDataFrame([
(1, "ivan@crm.ru", "", "N/A", "ACTIVE", 500.0),
(2, "", "anna@app.ru", "anna@old.ru", "UNKNOWN", None),
(3, None, None, None, "", float('nan')),
(4, "bob@crm.ru", "bob@app.ru", "bob@old.ru", "INACTIVE", 100.0),
], ["id", "email_crm", "email_app", "email_legacy", "status", "score"])
Нужно:
final_email- первый непустой email из трёх источников (nullif+coalesce)status_clean- NULL вместо""и"UNKNOWN"score_clean- NaN → NULL черезnanvlprofile_quality-"complete"если есть email и score,"partial"если только одно,"poor"если ничего
Решение
from pyspark.sql.functions import (
col, coalesce, nullif, nanvl, when, lit, trim, lower
)
result = (
customers
# 1. final_email: nullif убирает пустые строки
.withColumn(
"final_email",
coalesce(
nullif(trim(col("email_crm")), lit("")),
nullif(trim(col("email_app")), lit("")),
nullif(trim(col("email_legacy")), lit("N/A")),
)
)
# 2. status_clean: "" и "UNKNOWN" → NULL
.withColumn(
"status_clean",
nullif(nullif(col("status"), lit("")), lit("UNKNOWN"))
)
# 3. score_clean: NaN → NULL
.withColumn("score_clean", nanvl(col("score"), lit(None)))
# 4. profile_quality
.withColumn(
"profile_quality",
when(
col("final_email").isNotNull() & col("score_clean").isNotNull(),
"complete"
).when(
col("final_email").isNotNull() | col("score_clean").isNotNull(),
"partial"
).otherwise("poor")
)
.select("id", "final_email", "status_clean", "score_clean", "profile_quality")
)
result.show(truncate=False)
# +---+-----------+------------+-----------+---------------+
# | id|final_email|status_clean|score_clean|profile_quality|
# +---+-----------+------------+-----------+---------------+
# | 1|ivan@crm.ru| ACTIVE| 500.0| complete|
# | 2|anna@app.ru| null| null| partial|
# | 3| null| null| null| poor|
# | 4|bob@crm.ru | INACTIVE| 100.0| complete|
# +---+-----------+------------+-----------+---------------+
nullif(nullif(col("status"), lit("")), lit("UNKNOWN")) - цепочка двух nullif для нескольких sentinel-значений. Альтернатива - when(col("status").isin("", "UNKNOWN"), lit(None)).otherwise(col("status")).
Задание 2: Флаги активности и сегментация¶
Для DataFrame пользователей рассчитайте:
is_active- пользователь активен если: заходил за последние 30 дней И есть email И статус не"BLOCKED"risk_level-"high"если нет email или телефона,"medium"если есть только одно,"low"если оба естьpriority_score- число от 0 до 3: по 1 очку за email, за телефон, за вход за 7 дней
users = spark.createDataFrame([
(1, "ivan@mail.ru", "79001234567", 5, "ACTIVE"),
(2, "anna@mail.ru", None, 15, "ACTIVE"),
(3, None, "79009876543", 45, "ACTIVE"),
(4, None, None, 3, "BLOCKED"),
(5, "bob@mail.ru", "79005555555", None, "ACTIVE"),
], ["id", "email", "phone", "days_since_login", "status"])
Решение
from pyspark.sql.functions import col, when, lit
has_email = col("email").isNotNull()
has_phone = col("phone").isNotNull()
recent_login = col("days_since_login") <= 30
not_blocked = col("status") != "BLOCKED"
result = (
users
# 1. is_active: все три условия должны быть True
# NULL в days_since_login даёт NULL для <=30, что означает False в filter
.withColumn(
"is_active",
has_email & recent_login & not_blocked
)
# 2. risk_level: оба отсутствуют → high, одно → medium, оба есть → low
.withColumn(
"risk_level",
when(~has_email & ~has_phone, "high")
.when(has_email & has_phone, "low")
.otherwise("medium")
)
# 3. priority_score: сумма очков через when (условные 0/1)
.withColumn(
"priority_score",
when(has_email, 1).otherwise(0) +
when(has_phone, 1).otherwise(0) +
when(col("days_since_login") <= 7, 1).otherwise(0)
)
)
result.show()
# +---+--------+---------+------------------+--------+---------+----------+--------------+
# | id| email| phone|days_since_login | status|is_active|risk_level|priority_score|
# +---+--------+---------+------------------+--------+---------+----------+--------------+
# | 1|ivan@...|790012...| 5 | ACTIVE| true| low| 3|
# | 2|anna@...| null| 15| ACTIVE| true| medium| 2|
# | 3| null|790098...| 45| ACTIVE| false| medium| 1|
# | 4| null| null| 3| BLOCKED| false| high| 0|
# | 5|bob@... |790055...| null| ACTIVE| null| low| 2| ← NULL дней
# +---+--------+---------+------------------+--------+---------+----------+--------------+
days_since_login = NULL, поэтому NULL <= 30 = NULL, и is_active = email_ok AND NULL AND not_blocked = NULL (не false!). Если хотите явно false: coalesce(recent_login, lit(False)) перед использованием.
Задание 3: NULL-safe Change Detection¶
Есть два снапшота данных: current_snapshot и previous_snapshot. Нужно найти строки, где значение status изменилось, включая переходы NULL → значение и значение → NULL:
current = spark.createDataFrame([
(1, "ACTIVE"),
(2, "INACTIVE"),
(3, None),
(4, "ACTIVE"),
], ["id", "status"])
previous = spark.createDataFrame([
(1, "ACTIVE"), # не изменился
(2, None), # None → INACTIVE: изменился
(3, None), # None → None: не изменился
(4, "INACTIVE"), # INACTIVE → ACTIVE: изменился
], ["id", "status"])
Решение
from pyspark.sql.functions import col
# JOIN и сравнение: eqNullSafe находит «не изменился»
# ~ (NOT) инвертирует: находим строки где значение ИЗМЕНИЛОСЬ
changes = (
current.alias("cur")
.join(previous.alias("prev"), "id", "left")
.filter(
~col("cur.status").eqNullSafe(col("prev.status"))
)
.select(
col("id"),
col("prev.status").alias("old_status"),
col("cur.status").alias("new_status"),
)
)
changes.show()
# +---+----------+----------+
# | id|old_status|new_status|
# +---+----------+----------+
# | 2| null| INACTIVE| # null → INACTIVE: изменилось
# | 4| INACTIVE| ACTIVE| # INACTIVE → ACTIVE: изменилось
# +---+----------+----------+
# id=1: ACTIVE <=> ACTIVE → True → ~ → False → не изменился (корректно)
# id=3: None <=> None → True → ~ → False → не изменился (корректно)
eqNullSafe: NULL != "INACTIVE" → NULL → строка 2 не пройдёт фильтр и мы пропустим реальное изменение. eqNullSafe - единственный правильный способ сравнивать nullable колонки для CDC.
Задание 4: Conditional aggregation для KPI dashboard¶
Для датасета транзакций постройте агрегат с раздельными метриками по категориям в одном groupBy:
transactions = spark.createDataFrame([
("RU", "Electronics", 1500.0, "COMPLETED"),
("RU", "Clothing", 300.0, "COMPLETED"),
("RU", "Electronics", 200.0, "CANCELLED"),
("EU", "Electronics", 2000.0, "COMPLETED"),
("EU", "Clothing", 500.0, "PENDING"),
], ["region", "category", "amount", "status"])
Нужно по region: total_revenue, completed_revenue, electronics_revenue, cancelled_count, completion_rate (completed / total, округлить до 2 знаков).
Решение
from pyspark.sql.functions import (
col, when, sum as fsum, count, round as fround
)
result = (
transactions
.groupBy("region")
.agg(
fsum("amount").alias("total_revenue"),
# Выручка только COMPLETED заказов
fsum(
when(col("status") == "COMPLETED", col("amount")).otherwise(0)
).alias("completed_revenue"),
# Выручка только Electronics
fsum(
when(col("category") == "Electronics", col("amount")).otherwise(0)
).alias("electronics_revenue"),
# Количество CANCELLED (count(when) - считает только не-NULL)
count(
when(col("status") == "CANCELLED", 1)
).alias("cancelled_count"),
# Всего заказов для расчёта rate
count("*").alias("total_orders"),
count(
when(col("status") == "COMPLETED", 1)
).alias("completed_orders"),
)
.withColumn(
"completion_rate",
fround(col("completed_orders") / col("total_orders"), 2)
)
.drop("total_orders", "completed_orders")
)
result.show()
# +------+-------------+-----------------+-------------------+---------------+---------------+
# |region|total_revenue|completed_revenue|electronics_revenue|cancelled_count|completion_rate|
# +------+-------------+-----------------+-------------------+---------------+---------------+
# | RU| 2000.0| 1800.0| 1700.0| 1| 0.67|
# | EU| 2500.0| 2000.0| 2000.0| 0| 0.50|
# +------+-------------+-----------------+-------------------+---------------+---------------+
count(when(condition, 1)) - идиоматичный паттерн для conditional count: when без otherwise возвращает NULL при ложном условии, count не считает NULL. sum(when(condition, amount).otherwise(0)) - conditional sum с явным 0 для неподходящих строк.