Date/Time Functions: от строк до таймзон
Типы DateType и TimestampType, парсинг, арифметика дат, работа с Unix time и таймзонами, UTC-first архитектура и date partitioning в production lakehouse.
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¶
Дан лог входов пользователей. Рассчитайте:
first_login,last_login- первый и последний вход (через агрегацию)lifetime_days- дней между первым и последним входомdays_since_last_login- дней с последнего входа до сегодняis_churned-trueесли не заходил более 30 дней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|
# +--------+-------------------+--------------+--------+-----------------+
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) - нет.