Date/Time Functions: от строк до таймзон

Типы DateType и TimestampType, парсинг, арифметика дат, работа с Unix time и таймзонами, UTC-first архитектура и date partitioning в production lakehouse.

core

Date/Time Functions: от строк до таймзон

Дата и время - одна из самых коварных тем в Data Engineering. Источники присылают строки в десятках форматов, Kafka отдаёт миллисекунды от эпохи, веб-логи хранят UTC, а бизнес-аналитики требуют отчёты по московскому времени. При этом одна неверная настройка session.timeZone может незаметно сдвинуть все данные на несколько часов.

В этом уроке мы разберём всё: от внутреннего представления типов до правильной архитектуры хранения временны́х данных в lakehouse.

Внутреннее устройство: DateType и TimestampType

Прежде чем работать с датами, важно понять как Spark хранит их физически. Это объясняет поведение при конвертации, сравнении и работе с таймзонами.

Ключевые следствия этой модели:

  • DateType не имеет понятия о времени и часовом поясе - это просто число дней от эпохи. Два DateType из разных источников всегда сравниваются корректно.
  • TimestampType хранится как UTC-микросекунды, но отображается через spark.sql.session.timeZone. Это источник самых неочевидных багов: один и тот же Int64 выглядит по-разному в зависимости от настройки сессии.
  • Разница TimestampType - TimestampType даёт LongType (микросекунды), а DateType - DateType через datediff - IntegerType (дни).

Практическое правило: используйте DateType для всего, где время дня не важно (дата заказа, дата рождения, партиционная колонка). TimestampType - только для event-time с нужной точностью до секунды или меньше.

Парсинг: превращение строк в типы

to_date: строка → DateType

to_date(col, format) конвертирует строку в DateType. Формат задаётся в синтаксисе Java DateTimeFormatter. Если формат не указан - Spark пробует ISO 8601 (yyyy-MM-dd) и несколько других стандартных форматов:

from pyspark.sql.functions import to_date, col

df = spark.createDataFrame([
    ("2024-01-15",),
    ("15/01/2024",),
    ("Jan 15, 2024",),
    ("20240115",),
    (None,),
], ["raw_date"])

df.select(
    col("raw_date"),
    # ISO 8601 - format можно не указывать
    to_date(col("raw_date"), "yyyy-MM-dd").alias("iso"),
    # Европейский формат
    to_date(col("raw_date"), "dd/MM/yyyy").alias("eu"),
    # Текстовый месяц
    to_date(col("raw_date"), "MMM dd, yyyy").alias("text"),
    # Компактный числовой
    to_date(col("raw_date"), "yyyyMMdd").alias("compact"),
).show()
# +-----------+----------+----------+----------+----------+
# |   raw_date|       iso|        eu|      text|   compact|
# +-----------+----------+----------+----------+----------+
# | 2024-01-15|2024-01-15|      null|      null|      null|
# | 15/01/2024|      null|2024-01-15|      null|      null|
# |Jan 15, 2024|      null|      null|2024-01-15|      null|
# |   20240115|      null|      null|      null|2024-01-15|
# |       null|      null|      null|      null|      null|
# +-----------+----------+----------+----------+----------+

Важно: при несоответствии строки формату to_date возвращает NULL (не исключение). Это позволяет диагностировать проблемные записи - после to_date примените фильтр isNull() чтобы найти строки, которые не распарсились.

to_timestamp: строка → TimestampType

to_timestamp(col, format) работает аналогично, но возвращает TimestampType. Формат должен включать компоненты времени. Результат интерпретируется как локальное время в spark.sql.session.timeZone и конвертируется в UTC-микросекунды для хранения:

from pyspark.sql.functions import to_timestamp

df = spark.createDataFrame([
    ("2024-01-15 12:30:00",),
    ("2024-01-15T12:30:00Z",),          # ISO 8601 с Z (UTC)
    ("15-01-2024 12:30:00.500",),       # с миллисекундами
    ("2024-01-15 12:30:00+03:00",),     # с offset
], ["raw_ts"])

df.select(
    col("raw_ts"),
    to_timestamp(col("raw_ts"), "yyyy-MM-dd HH:mm:ss").alias("ts_basic"),
    to_timestamp(col("raw_ts"), "yyyy-MM-dd'T'HH:mm:ssX").alias("ts_iso"),
    to_timestamp(col("raw_ts"), "dd-MM-yyyy HH:mm:ss.SSS").alias("ts_ms"),
).show(truncate=False)

Символ 'T' в паттерне - это литерал (кавычки обязательны, иначе T интерпретируется как timezone). X - timezone offset в формате +HH, +HHMM или Z для UTC. SSS - миллисекунды.

date_format: дата → строка

date_format(col, pattern) - обратная операция: типизированную дату или timestamp превращает обратно в строку по заданному паттерну. Применяется при выгрузке данных в системы с фиксированным форматом:

from pyspark.sql.functions import date_format, current_timestamp

spark.range(1).select(
    current_timestamp().alias("ts"),
    date_format(current_timestamp(), "yyyy-MM-dd").alias("date_only"),
    date_format(current_timestamp(), "dd.MM.yyyy HH:mm").alias("ru_format"),
    date_format(current_timestamp(), "yyyyMMdd_HHmmss").alias("file_suffix"),
    date_format(current_timestamp(), "EEEE, d MMMM yyyy").alias("human"),
).show(truncate=False)
# +-------------------+----------+----------------+-----------------+------------------------+
# |ts                 |date_only |ru_format       |file_suffix      |human                   |
# +-------------------+----------+----------------+-----------------+------------------------+
# |2024-01-15 12:30:00|2024-01-15|15.01.2024 12:30|20240115_123000  |Monday, 15 January 2024 |
# +-------------------+----------+----------------+-----------------+------------------------+

Форматные паттерны Java DateTimeFormatter

Spark использует синтаксис java.time.format.DateTimeFormatter. Отличается от Python strftime:

Паттерн Значение Пример
yyyy 4-значный год 2024
yy 2-значный год 24
MM месяц с ведущим нулём 01-12
M месяц без ведущего нуля 1-12
MMM сокращённое название Jan, Feb
MMMM полное название January
dd день с ведущим нулём 01-31
d день без ведущего нуля 1-31
HH час 24h с нулём 00-23
hh час 12h с нулём 01-12
mm минуты 00-59
ss секунды 00-59
SSS миллисекунды 000-999
EEEE полное название дня недели Monday
E сокращённое название дня Mon
X timezone offset +03, Z
z timezone abbreviation UTC, MSK

Отличие от Python: в Python %m - месяц, %d - день. В Spark MM - месяц, dd - день. Строчные mm в Spark - это минуты, а не месяц.

Извлечение компонентов даты

После приведения строки к DateType или TimestampType можно извлекать отдельные части через специализированные функции:

from pyspark.sql.functions import (
    year, month, dayofmonth, dayofweek, dayofyear,
    quarter, weekofyear, hour, minute, second,
    to_date, to_timestamp, col
)

df = spark.createDataFrame([("2024-03-15 14:35:22",)], ["ts_str"])
df = df.withColumn("ts", to_timestamp(col("ts_str"), "yyyy-MM-dd HH:mm:ss")) \
       .withColumn("dt", to_date(col("ts_str"), "yyyy-MM-dd HH:mm:ss"))

df.select(
    # Компоненты даты
    year("dt").alias("year"),          # 2024
    month("dt").alias("month"),        # 3
    dayofmonth("dt").alias("day"),     # 15
    quarter("dt").alias("quarter"),    # 1
    weekofyear("dt").alias("week"),    # 11
    dayofweek("dt").alias("dow"),      # 6 (1=Sun, 2=Mon, ..., 7=Sat)
    dayofyear("dt").alias("doy"),      # 75

    # Компоненты времени (только из TimestampType)
    hour("ts").alias("hour"),          # 14
    minute("ts").alias("minute"),      # 35
    second("ts").alias("second"),      # 22
).show()

Нюанс dayofweek: Spark использует нумерацию SQL: 1 = воскресенье, 2 = понедельник, ..., 7 = суббота. Это отличается от Python weekday() (0 = понедельник). Для определения выходного дня в российском контексте: dayofweek(col("dt")).isin(1, 7) - воскресенье и суббота.

Функции работают как с DateType, так и с TimestampType. Функции hour, minute, second на DateType вернут 0.

Арифметика дат

date_add и date_sub: сдвиг на дни

date_add(col, n) прибавляет n дней к дате. date_sub(col, n) - вычитает. Учитывают реальную длину месяца и не дают неправильных дат вроде 30 февраля:

from pyspark.sql.functions import date_add, date_sub, col, to_date, lit

df = spark.createDataFrame([("2024-01-28",), ("2024-02-28",)], ["dt_str"])
df = df.withColumn("dt", to_date("dt_str", "yyyy-MM-dd"))

df.select(
    col("dt"),
    date_add(col("dt"), 5).alias("plus_5_days"),   # +5 дней
    date_sub(col("dt"), 7).alias("minus_7_days"),  # -7 дней
    # n может быть колонкой, не только литералом
    date_add(col("dt"), col("offset")).alias("dynamic") if "offset" in df.columns else lit(None),
).show()
# +----------+----------+----------+
# |        dt|plus_5_days|minus_7_days|
# +----------+----------+----------+
# |2024-01-28|2024-02-02|2024-01-21|  # корректно пересёк границу месяца
# |2024-02-28|2024-03-04|2024-02-21|  # корректно пересёк в марте
# +----------+----------+----------+

Типичное применение: вычисление окна ретеншн (date_add(first_login, 30)) или инкрементальной обработки (date_sub(current_date(), 1) - вчерашние данные).

datediff: разница в днях

datediff(end, start) возвращает IntegerType - количество дней между двумя датами. Отрицательное значение если end < start:

from pyspark.sql.functions import datediff, current_date

df = spark.createDataFrame([
    (1, "2024-01-01", "2024-03-15"),
    (2, "2023-12-01", "2024-01-01"),
    (3, "2024-06-01", "2024-01-01"),  # отрицательное
], ["user_id", "first_login", "last_login"])

df.withColumn("first_dt",    to_date("first_login",  "yyyy-MM-dd")) \
  .withColumn("last_dt",     to_date("last_login",   "yyyy-MM-dd")) \
  .select(
    col("user_id"),
    # Время жизни пользователя (дни между первым и последним входом)
    datediff(col("last_dt"),    col("first_dt")).alias("lifetime_days"),
    # Количество дней с регистрации до сегодня
    datediff(current_date(),    col("first_dt")).alias("days_since_signup"),
).show()
# +-------+-------------+-----------------+
# |user_id|lifetime_days|days_since_signup|
# +-------+-------------+-----------------+
# |      1|           74|              ...| 
# |      2|           31|              ...|
# |      3|         -152|              ...|
# +-------+-------------+-----------------+

add_months и months_between: месячная арифметика

add_months(col, n) добавляет n целых месяцев с корректной обработкой «конца месяца»: если исходная дата - последний день месяца, результат тоже последний день целевого месяца:

from pyspark.sql.functions import add_months, months_between

df = spark.createDataFrame([
    ("2024-01-31",),  # последний день января
    ("2024-01-15",),
], ["dt_str"])
df = df.withColumn("dt", to_date("dt_str", "yyyy-MM-dd"))

df.select(
    col("dt"),
    add_months(col("dt"), 1).alias("plus_1m"),   # 2024-01-31 → 2024-02-29 (не 03-02!)
    add_months(col("dt"), -1).alias("minus_1m"),  # вычитание
    add_months(col("dt"), 12).alias("plus_1y"),   # +1 год через месяцы
).show()
# +----------+----------+----------+----------+
# |        dt|  plus_1m |  minus_1m|  plus_1y |
# +----------+----------+----------+----------+
# |2024-01-31|2024-02-29|2023-12-31|2025-01-31|  # 29 фев - последний день февраля 2024!
# |2024-01-15|2024-02-15|2023-12-15|2025-01-15|
# +----------+----------+----------+----------+

months_between(end, start) - дробное число месяцев между датами (учитывает дни):

# Расчёт tenure в месяцах для subscription analytics
df.withColumn("tenure_months",
    months_between(current_date(), col("dt")).cast("int")
)
# 2024-01-15 → 2024-03-15: tenure_months = 2

last_day: последний день месяца

last_day(col) возвращает последний день месяца для данной даты. Используется в финансовой аналитике для расчёта конца отчётного периода:

from pyspark.sql.functions import last_day

df.select(
    col("dt"),
    last_day(col("dt")).alias("month_end"),
    # Количество дней до конца месяца
    datediff(last_day(col("dt")), col("dt")).alias("days_to_month_end"),
)
# 2024-01-15 → month_end=2024-01-31, days_to_month_end=16
# 2024-02-01 → month_end=2024-02-29, days_to_month_end=28

trunc и date_trunc: округление вниз

trunc(col, unit) округляет DateType до начала указанного периода. date_trunc(unit, col) делает то же для TimestampType с единицами времени вплоть до секунды:

from pyspark.sql.functions import trunc, date_trunc

df = spark.createDataFrame([("2024-03-15 14:35:22",)], ["ts_str"])
df = df.withColumn("ts", to_timestamp("ts_str", "yyyy-MM-dd HH:mm:ss")) \
       .withColumn("dt", to_date("ts_str",      "yyyy-MM-dd HH:mm:ss"))

df.select(
    # trunc для DateType
    trunc(col("dt"), "month").alias("month_start"),  # 2024-03-01
    trunc(col("dt"), "year").alias("year_start"),    # 2024-01-01
    trunc(col("dt"), "quarter").alias("qtr_start"),  # 2024-01-01 (Q1)

    # date_trunc для TimestampType
    date_trunc("hour",  col("ts")).alias("hour_trunc"),    # 2024-03-15 14:00:00
    date_trunc("day",   col("ts")).alias("day_trunc"),     # 2024-03-15 00:00:00
    date_trunc("month", col("ts")).alias("month_trunc"),   # 2024-03-01 00:00:00
    date_trunc("week",  col("ts")).alias("week_trunc"),    # 2024-03-11 00:00:00 (пн)
).show(truncate=False)

trunc и date_trunc - фундамент для построения временны́х агрегатов: groupBy(trunc(col("dt"), "month")) даёт ежемесячные итоги без ручного извлечения year/month. date_trunc("hour", col("ts")) - основа почасовых агрегатов в streaming pipeline.

current_date и current_timestamp

from pyspark.sql.functions import current_date, current_timestamp

df.withColumn("today",      current_date())       # DateType, дата выполнения джобы
  .withColumn("now",        current_timestamp())  # TimestampType, с точностью до µs
  .withColumn("ingested_at", current_timestamp()) # audit-колонка

current_date() и current_timestamp() вычисляются один раз при начале выполнения Action (не при построении плана). Все строки в одном датасете получат одинаковое значение. Это важно для audit-колонок: ingested_at будет одинаковым для всего batch.

Unix Time: мост между системами

Unix timestamp (epoch seconds) - универсальный формат обмена временем между системами: Kafka события, веб-логи nginx, API-ответы. Spark умеет работать с ним нативно.

unix_timestamp: дата → epoch seconds

unix_timestamp(col, format) конвертирует строку или DateType/TimestampType в количество секунд от эпохи (Int64). Без аргументов - возвращает текущее время:

from pyspark.sql.functions import unix_timestamp

df.select(
    col("ts_str"),
    unix_timestamp(col("ts_str"), "yyyy-MM-dd HH:mm:ss").alias("epoch_sec"),
    # Из TimestampType без формата
    unix_timestamp(col("ts")).alias("epoch_from_ts"),
)
# 2024-01-15 12:30:00 UTC → 1705319400

from_unixtime: epoch seconds → строка

from_unixtime(col, format) конвертирует epoch seconds обратно в строку заданного формата:

from pyspark.sql.functions import from_unixtime

df = spark.createDataFrame([(1705319400,), (1705320000,)], ["epoch_sec"])

df.select(
    col("epoch_sec"),
    from_unixtime(col("epoch_sec"), "yyyy-MM-dd HH:mm:ss").alias("datetime_str"),
    # Без формата - ISO 8601 по умолчанию
    from_unixtime(col("epoch_sec")).alias("datetime_default"),
    # Конвертировать в TimestampType
    to_timestamp(from_unixtime(col("epoch_sec"))).alias("ts"),
).show()

Ошибка с миллисекундами

Самая частая ошибка при работе с epoch: Kafka, Java System.currentTimeMillis() и большинство современных систем отдают миллисекунды, а from_unixtime ожидает секунды:

# ОШИБКА: kafka_ts в миллисекундах
# from_unixtime(col("kafka_ts")) → дата в 53 000 году!

# ПРАВИЛЬНО: делить на 1000
from_unixtime((col("kafka_ts") / 1000).cast("long"), "yyyy-MM-dd HH:mm:ss")

# Или через cast TimestampType напрямую (Spark 3.0+):
# (col("kafka_ts") / 1000).cast("timestamp")  # делит на 1000 и конвертирует

Таймзоны: главная боль инженера

Таймзоны - источник самых трудноуловимых багов в data pipeline. Одна неправильная настройка может незаметно сдвинуть все события на несколько часов, что искажает метрики, когортный анализ и любую аналитику по времени суток.

Все источники фиксируют один и тот же момент времени (12:00 UTC), но записывают его в своей локальной зоне. При ingestion нужно привести всё к UTC. При построении витрин - конвертировать обратно в нужную зону для читателя.

spark.sql.session.timeZone

Это системная настройка, определяющая как Spark интерпретирует строки без явного timezone и как отображает TimestampType:

# Проверить текущую настройку
spark.conf.get("spark.sql.session.timeZone")  # обычно "UTC" или системная TZ сервера

# Изменить для сессии
spark.conf.set("spark.sql.session.timeZone", "Europe/Moscow")

# Опасный пример: один и тот же timestamp выглядит по-разному
spark.conf.set("spark.sql.session.timeZone", "UTC")
df.select(to_timestamp(lit("2024-01-15 12:00:00"), "yyyy-MM-dd HH:mm:ss")).show()
# 2024-01-15 12:00:00  ← интерпретируется как UTC

spark.conf.set("spark.sql.session.timeZone", "Europe/Moscow")
df.select(to_timestamp(lit("2024-01-15 12:00:00"), "yyyy-MM-dd HH:mm:ss")).show()
# 2024-01-15 12:00:00  ← отображается, но внутри хранится 09:00 UTC!

Правило: на production кластере явно устанавливайте spark.sql.session.timeZone = UTC в конфигурации Spark. Никогда не полагайтесь на системную timezone сервера - она может отличаться на разных нодах кластера.

from_utc_timestamp и to_utc_timestamp

to_utc_timestamp(col, tz) - конвертирует локальное время заданной TZ в UTC. Применяется при ingestion:

from pyspark.sql.functions import to_utc_timestamp, from_utc_timestamp

df = spark.createDataFrame([
    ("2024-03-15 15:00:00", "Europe/Moscow"),
    ("2024-03-15 07:00:00", "America/New_York"),
], ["local_ts_str", "timezone"])

df.select(
    col("local_ts_str"),
    col("timezone"),
    # Конвертация в UTC - статический TZ
    to_utc_timestamp(
        to_timestamp(col("local_ts_str"), "yyyy-MM-dd HH:mm:ss"),
        "Europe/Moscow"
    ).alias("utc_ts"),
    # Динамический TZ из колонки
    to_utc_timestamp(
        to_timestamp(col("local_ts_str"), "yyyy-MM-dd HH:mm:ss"),
        col("timezone")
    ).alias("utc_from_col"),
).show(truncate=False)
# Оба варианта дают 2024-03-15 12:00:00 UTC для московского источника

from_utc_timestamp(col, tz) - обратное преобразование: из UTC в локальную TZ для витрин:

# Данные хранятся в UTC, пользователь хочет видеть московское время
df.select(
    from_utc_timestamp(col("event_ts_utc"), "Europe/Moscow").alias("event_ts_msk")
)

DST: летнее время ломает аналитику

Daylight Saving Time (переход летнего/зимнего времени) - особый источник проблем. В день перехода часы переводятся: время 02:30 в Europe/London может не существовать или существовать дважды:

# Переход на летнее время в Европе: 31 марта 2024 в 02:00 → 03:00
# В этот момент 02:30 Europe/London не существует!

to_utc_timestamp(
    to_timestamp(lit("2024-03-31 02:30:00"), "yyyy-MM-dd HH:mm:ss"),
    "Europe/London"
)
# Spark решает неоднозначность самостоятельно (обычно pre-transition)
# Лучшее решение: хранить события уже в UTC с источника

Практическое правило: если ваш источник данных (мобильное приложение, веб-сервис) может писать времена в разных TZ или при переходе DST - настройте источник на отдачу UTC или epoch milliseconds. Любой постфактум-пересчёт TZ на стороне Spark будет иметь edge cases.

UTC-first: архитектурный стандарт

Современный production lakehouse всегда хранит timestamps в UTC. Это не просто рекомендация - это стандарт, от которого зависит корректность всей downstream аналитики:

Слой Формат Обоснование
Bronze как пришло из источника сохраняем raw data без изменений
Silver TimestampType в UTC нормализация при ingestion
Gold / витрины UTC или явная TZ-конвертация зависит от аудитории отчёта

Date Partitioning и Predicate Pushdown

Партиционирование по дате - стандарт для хранения событийных данных в lakehouse. Правильное партиционирование и правильные фильтры дают partition pruning - Spark пропускает целые папки с данными:

Ключевой принцип: применение любой функции к партиционной колонке ломает partition pruning. Spark не может заглянуть внутрь функции и понять, какие партиции затронуты.

# ПРАВИЛЬНО: партиционная колонка - отдельная, строковая
silver.write.partitionBy("p_date").parquet("/data/silver/events/")

# При чтении фильтровать по строке p_date, не по derived timestamp
spark.read.parquet("/data/silver/events/") \
    .filter(col("p_date") == "2024-01-15")  # ← partition pruning работает

# АНТИПАТТЕРН: фильтр через функцию не даёт pruning
spark.read.parquet("/data/silver/events/") \
    .filter(year(col("event_ts")) == 2024)  # ← читает ВСЕ партиции!

Партиционная колонка правильно

from pyspark.sql.functions import date_format, to_date

# При записи: создать p_date как строку YYYY-MM-DD из timestamp
silver = raw.withColumn(
    "p_date",
    date_format(to_timestamp(col("event_ts"), "yyyy-MM-dd HH:mm:ss"), "yyyy-MM-dd")
)

silver.write \
    .mode("append") \
    .partitionBy("p_date") \
    .parquet("/data/silver/events/")

# При чтении: фильтр по строке - мгновенный partition scan
events_jan = spark.read.parquet("/data/silver/events/") \
    .filter(col("p_date").between("2024-01-01", "2024-01-31"))
    # between на строках работает как лексикографическое сравнение,
    # что совпадает с хронологическим порядком для формата YYYY-MM-DD

Формат yyyy-MM-dd для p_date критически важен: только он гарантирует что лексикографический порядок строк совпадает с хронологическим порядком дат. dd-MM-yyyy или MM/dd/yyyy сломают сортировку.

Event time vs processing time

В потоковой обработке различают два понятия времени, которые важно не путать:

Event time - момент когда событие произошло в реальном мире. Записывается источником данных. Может прийти с опозданием (late arrival).

Processing time - момент когда Spark обработал событие. Это current_timestamp() в момент выполнения джобы.

from pyspark.sql.functions import current_timestamp, to_timestamp, col

df.select(
    col("event_id"),
    to_timestamp(col("event_ts_str"), "yyyy-MM-dd HH:mm:ss").alias("event_time"),    # когда произошло
    current_timestamp().alias("processed_at"),                                        # когда обработали
    # Задержка обработки (processing lag)
    (unix_timestamp(current_timestamp()) -
     unix_timestamp(to_timestamp(col("event_ts_str"), "yyyy-MM-dd HH:mm:ss"))
    ).alias("lag_seconds"),
)

Для batch pipeline партиционирование почти всегда должно быть по event time, а не по processing time. Иначе поздно пришедшие события попадают в «неправильную» партицию относительно момента события.

Anti-patterns в datetime processing

# 1. Строки вместо нативных типов - ломает арифметику и partition pruning
df.withColumn("date_str", lit("2024-01-15"))  # плохо
df.withColumn("date_col", to_date(lit("2024-01-15")))  # хорошо

# 2. Python datetime через UDF - pickle overhead + Catalyst не оптимизирует
from datetime import datetime
@udf("string")
def format_date_udf(s):
    return datetime.strptime(s, "%Y-%m-%d").strftime("%d.%m.%Y")
# Используй date_format(to_date(col, "yyyy-MM-dd"), "dd.MM.yyyy") - в 30× быстрее

# 3. Игнорирование timezone при to_timestamp - непредсказуемое поведение
to_timestamp(col("ts_str"))  # зависит от session.timeZone кластера

# 4. Применение функций к партиционной колонке - нет pruning
.filter(year(col("p_date")) == 2024)  # читает ВСЕ партиции
.filter(col("p_date").startswith("2024"))  # partition pruning работает!

# 5. epoch milliseconds → from_unixtime без /1000
from_unixtime(col("kafka_ts"))              # неверно (год ~53000)
from_unixtime(col("kafka_ts") / 1000)      # верно

Практика

Задание 1: Парсинг разнородных форматов

Дан DataFrame с датами в разных форматах из разных источников. Привести все к единому DateType:

df = spark.createDataFrame([
    (1, "2024-01-15"),           # ISO 8601
    (2, "15/01/2024"),           # европейский
    (3, "Jan 15, 2024"),         # текстовый месяц
    (4, "20240115"),             # компактный
    (5, "2024.01.15"),           # точки
    (6, "bad-date"),             # невалидный
    (7, None),                   # null
], ["id", "raw_date"])

Нужно: создать колонку parsed_date (DateType), колонку parse_error (BooleanType) - true для невалидных или null строк, и вывести количество ошибок по источнику.

Решение

from pyspark.sql.functions import (
    to_date, coalesce, col, when
)

# Пробуем форматы по очереди через coalesce - вернёт первый не-null
result = df.withColumn(
    "parsed_date",
    coalesce(
        to_date(col("raw_date"), "yyyy-MM-dd"),
        to_date(col("raw_date"), "dd/MM/yyyy"),
        to_date(col("raw_date"), "MMM dd, yyyy"),
        to_date(col("raw_date"), "yyyyMMdd"),
        to_date(col("raw_date"), "yyyy.MM.dd"),
    )
).withColumn(
    "parse_error",
    col("parsed_date").isNull()  # null если ни один формат не подошёл
)

result.show()
# +---+------------+----------+-----------+
# | id|    raw_date|parsed_date|parse_error|
# +---+------------+----------+-----------+
# |  1|  2024-01-15|2024-01-15|      false|
# |  2|  15/01/2024|2024-01-15|      false|
# |  3|Jan 15, 2024|2024-01-15|      false|
# |  4|    20240115|2024-01-15|      false|
# |  5|  2024.01.15|2024-01-15|      false|
# |  6|    bad-date|      null|       true|
# |  7|        null|      null|       true|
# +---+------------+----------+-----------+

result.filter(col("parse_error")).count()  # 2 ошибки
coalesce пробует выражения слева направо и возвращает первое не-null. to_date при несовпадении формата возвращает null (не исключение), поэтому цепочка безопасна. На больших датасетах этот подход эффективен - Catalyst объединяет все to_date в один проход.

Задание 2: User retention metrics

Дан лог входов пользователей. Рассчитайте:

  1. first_login, last_login - первый и последний вход (через агрегацию)
  2. lifetime_days - дней между первым и последним входом
  3. days_since_last_login - дней с последнего входа до сегодня
  4. is_churned - true если не заходил более 30 дней
  5. cohort_month - месяц первого входа (начало месяца через trunc)
events = spark.createDataFrame([
    (1, "2024-01-05"),
    (1, "2024-01-20"),
    (1, "2024-02-10"),
    (2, "2024-01-10"),
    (2, "2024-01-15"),
    (3, "2023-12-01"),
], ["user_id", "login_date_str"])
Решение

from pyspark.sql.functions import (
    to_date, min as fmin, max as fmax,
    datediff, current_date, trunc, col
)

events = events.withColumn("login_date", to_date("login_date_str", "yyyy-MM-dd"))

user_metrics = (
    events
    .groupBy("user_id")
    .agg(
        fmin("login_date").alias("first_login"),
        fmax("login_date").alias("last_login"),
    )
    .withColumn("lifetime_days",
        datediff(col("last_login"), col("first_login"))
    )
    .withColumn("days_since_last_login",
        datediff(current_date(), col("last_login"))
    )
    .withColumn("is_churned",
        datediff(current_date(), col("last_login")) > 30
    )
    .withColumn("cohort_month",
        trunc(col("first_login"), "month")  # 2024-01-05 → 2024-01-01
    )
)

user_metrics.show()
# +-------+-----------+----------+-------------+---------------------+----------+-----------+
# |user_id|first_login|last_login|lifetime_days|days_since_last_login|is_churned|cohort_month|
# +-------+-----------+----------+-------------+---------------------+----------+-----------+
# |      1| 2024-01-05|2024-02-10|           36|                 ...|      ...|  2024-01-01|
# |      2| 2024-01-10|2024-01-15|            5|                 ...|      ...|  2024-01-01|
# |      3| 2023-12-01|2023-12-01|            0|                 ...|     true|  2023-12-01|
# +-------+-----------+----------+-------------+---------------------+----------+-----------+
trunc(col("first_login"), "month") даёт первый день месяца - это стандартный способ определить когортный месяц. Все пользователи, зарегистрировавшиеся в январе 2024, получат cohort_month = 2024-01-01, независимо от точной даты.

Задание 3: Timezone normalization pipeline

Дан DataFrame с событиями из трёх регионов. Каждое событие содержит локальное время и TZ. Нужно привести к UTC и добавить аналитические измерения:

events = spark.createDataFrame([
    (1, "2024-03-15 15:00:00", "Europe/Moscow"),
    (2, "2024-03-15 07:00:00", "America/New_York"),
    (3, "2024-03-15 20:00:00", "Asia/Yekaterinburg"),
    (4, "2024-03-15 12:00:00", "UTC"),
], ["event_id", "local_ts_str", "timezone"])

Нужно: event_ts_utc (TimestampType UTC), event_date_utc (DateType), hour_utc (Int), is_business_hours - true если час UTC в 9-18.

Решение

from pyspark.sql.functions import (
    to_timestamp, to_utc_timestamp, to_date,
    hour, col, when
)

result = (
    events
    # 1. Парсим строку в TimestampType (без TZ - считается как local)
    .withColumn("local_ts",
        to_timestamp(col("local_ts_str"), "yyyy-MM-dd HH:mm:ss")
    )
    # 2. Конвертируем в UTC с учётом TZ из колонки
    .withColumn("event_ts_utc",
        to_utc_timestamp(col("local_ts"), col("timezone"))
    )
    # 3. Извлекаем дату и час в UTC
    .withColumn("event_date_utc", to_date(col("event_ts_utc")))
    .withColumn("hour_utc",       hour(col("event_ts_utc")))
    # 4. Бизнес-часы по UTC (9:00-18:00)
    .withColumn("is_business_hours",
        col("hour_utc").between(9, 17)
    )
    .select("event_id", "event_ts_utc", "event_date_utc",
            "hour_utc", "is_business_hours")
)

result.show(truncate=False)
# +--------+-------------------+--------------+--------+-----------------+
# |event_id|event_ts_utc       |event_date_utc|hour_utc|is_business_hours|
# +--------+-------------------+--------------+--------+-----------------+
# |       1|2024-03-15 12:00:00|    2024-03-15|      12|             true|
# |       2|2024-03-15 12:00:00|    2024-03-15|      12|             true|
# |       3|2024-03-15 15:00:00|    2024-03-15|      15|             true|
# |       4|2024-03-15 12:00:00|    2024-03-15|      12|             true|
# +--------+-------------------+--------------+--------+-----------------+
Все четыре события - разные локальные времена, но три из них (MSK 15:00, EST 07:00, UTC 12:00) соответствуют одному UTC-моменту. to_utc_timestamp(col, col("timezone")) принимает TZ из колонки - это позволяет обрабатывать разнородные источники в одном DataFrame без UDF.

Задание 4: Построение аналитических измерений

Для датасета событий создайте полный набор временны́х измерений для BI-витрины:

events = spark.createDataFrame([
    (1, "2024-01-15 09:30:00"),
    (2, "2024-07-04 14:15:00"),
    (3, "2024-12-31 23:59:59"),
], ["event_id", "event_ts_str"])

Нужно добавить: p_date (строка для партиции), event_date, year, month, day, quarter, week_of_year, day_of_week_name, is_weekend, hour, day_part (morning/afternoon/evening/night).

Решение

from pyspark.sql.functions import (
    to_timestamp, to_date, date_format,
    year, month, dayofmonth, quarter,
    weekofyear, dayofweek, hour,
    when, col
)

result = (
    events
    .withColumn("ts", to_timestamp("event_ts_str", "yyyy-MM-dd HH:mm:ss"))
    # Партиционная колонка - строка для partition pruning
    .withColumn("p_date",         date_format(col("ts"), "yyyy-MM-dd"))
    .withColumn("event_date",     to_date(col("ts")))
    # Временны́е измерения
    .withColumn("year",           year(col("ts")))
    .withColumn("month",          month(col("ts")))
    .withColumn("day",            dayofmonth(col("ts")))
    .withColumn("quarter",        quarter(col("ts")))
    .withColumn("week_of_year",   weekofyear(col("ts")))
    # day_of_week: 1=Sun, 7=Sat → текстовое название
    .withColumn("day_of_week_name", date_format(col("ts"), "EEEE"))
    # Выходной: воскресенье (1) или суббота (7)
    .withColumn("is_weekend",     dayofweek(col("ts")).isin(1, 7))
    .withColumn("hour",           hour(col("ts")))
    # Часть дня
    .withColumn("day_part",
        when(col("hour").between(6, 11),  "morning")
        .when(col("hour").between(12, 17), "afternoon")
        .when(col("hour").between(18, 22), "evening")
        .otherwise("night")
    )
    .drop("ts")
)

result.show(truncate=False)
# +--------+----------+----------+----+-----+---+-------+------------+----------+-------------+----------+----+---------+
# |event_id|p_date    |event_date|year|month|day|quarter|week_of_year|day_of_week|is_weekend|hour|day_part |
# +--------+----------+----------+----+-----+---+-------+------------+----------+----------+----+---------+
# |       1|2024-01-15|2024-01-15|2024|    1| 15|      1|           3|    Monday|     false|   9| morning |
# |       2|2024-07-04|2024-07-04|2024|    7|  4|      3|          27|  Thursday|     false|  14|afternoon|
# |       3|2024-12-31|2024-12-31|2024|   12| 31|      4|           1|   Tuesday|     false|  23|   night |
# +--------+----------+----------+----+-----+---+-------+------------+----------+----------+----+---------+
p_date как строка формата yyyy-MM-dd - именно такой тип нужно использовать при partitionBy("p_date"). Фильтрация .filter(col("p_date").between("2024-01-01", "2024-01-31")) даст partition pruning. Фильтрация .filter(year(col("event_date")) == 2024) - нет.