Conditional Columns: when/otherwise, coalesce, nullif и NULL-безопасные сравнения

Условная логика в DataFrame API: CASE WHEN, coalesce, nullif, NaN vs NULL, NULL-safe equality и трёхзначная логика SQL в production ETL.

core

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 заражает всё», потому что результат известен независимо от NULL
  • TRUE 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"])

Нужно:

  1. final_email - первый непустой email из трёх источников (nullif + coalesce)
  2. status_clean - NULL вместо "" и "UNKNOWN"
  3. score_clean - NaN → NULL через nanvl
  4. profile_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 пользователей рассчитайте:

  1. is_active - пользователь активен если: заходил за последние 30 дней И есть email И статус не "BLOCKED"
  2. risk_level - "high" если нет email или телефона, "medium" если есть только одно, "low" если оба есть
  3. 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 дней
# +---+--------+---------+------------------+--------+---------+----------+--------------+
Обратите внимание на строку 5: 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 для неподходящих строк.