String Functions: Spark SQL против Python UDF
Встроенные строковые функции PySpark - regexp_replace, concat_ws, split, substring и их преимущества над Python UDF в production ETL.
String Functions: Spark SQL против Python UDF¶
Строковые данные - самый «грязный» вид данных в pipeline. Имена клиентов с опечатками, телефоны в восьми разных форматах, логи сервисов как raw text, JSON-в-строке, CSV-в-строке. Data Engineer тратит огромную долю времени на нормализацию текстовых полей прежде чем они станут пригодны для аналитики.
В этом уроке мы разберём, почему встроенные строковые функции Spark - первый инструмент, который нужно освоить, прежде чем писать первый Python UDF, изучим полный арсенал функций и применим его к реальным задачам нормализации данных.
Почему UDF - крайний случай¶
Когда начинающие Data Engineers сталкиваются со строковой обработкой, первый инстинкт - написать Python-функцию. Python-разработчики привыкли к богатой стандартной библиотеке re, str. Но в контексте Spark это антипаттерн с серьёзными последствиями для производительности.
Python Barrier: путь данных через границу JVM¶
Spark Runtime работает на JVM (Java Virtual Machine). Executor хранит данные в формате UnsafeRow (off-heap бинарный формат) и обрабатывает строки JVM-байткодом. Когда вы вызываете Python UDF - каждая строка пересекает границу JVM:
Что происходит для каждой строки датасета:
- JVM сериализует UnsafeRow в Python-объект через Pickle - дорогостоящая операция, требующая копирования данных
- Данные отправляются в отдельный Python-процесс (он запускается рядом с JVM-воркером)
- Python выполняет вашу функцию и формирует результат
- Результат сериализуется обратно в байтовый поток (Pickle)
- JVM десериализует и записывает результат в новый UnsafeRow
Накладные расходы: 2 сериализации на строку, переключение контекста между процессами, невозможность векторизации операций. На 100 млн строк это занимает 4 минуты против 8 секунд у встроенных функций - разница в 30×.
Catalyst не видит внутренности UDF¶
Catalyst Optimizer - сердце производительности Spark. Он переписывает логический план: убирает лишние колонки, pushdown-ит фильтры в источник данных, объединяет несколько шагов в один. Но @udf-функции для Catalyst полностью непрозрачны:
Когда Spark встречает regexp_replace() - он включает её в WholeStageCodeGen: компилирует все трансформации в один JVM-байткод-проход без накладных расходов вызовов функций. Когда встречает Python UDF - ставит барьер BatchEvalPython и выключает весь CodeGen для этой части плана.
Vectorized UDF (Pandas UDF): лучше, но не замена¶
Pandas UDF работает с батчами строк через Apache Arrow без пословной сериализации. Это быстрее обычного UDF в 5-10×, но всё равно пересекает JVM-Python границу батчами, а Catalyst не оптимизирует содержимое функции.
Правило: Если задача решается встроенной функцией - используй её. Pandas UDF - когда нужна ML-библиотека или нестандартная логика. Обычный Python UDF - последний resort при отсутствии альтернатив.
Базовый инвентарь строковых функций¶
concat и concat_ws: конкатенация строк¶
concat() объединяет любое количество строк. Но у неё критичное поведение: если хотя бы один аргумент равен NULL - результат тоже NULL. concat_ws() (concat with separator) решает эту проблему - она игнорирует NULL и вставляет разделитель только между непустыми значениями:
from pyspark.sql.functions import concat, concat_ws, lit, col
df = spark.createDataFrame([
("Иван", "Петров", None),
("Анна", None, "Смирнова"),
], ["first", "last", "maiden"])
df.select(
# concat: NULL заражает весь результат
concat(col("first"), lit(" "), col("last")).alias("full_name_unsafe"),
# concat_ws: NULL-значения пропускаются
concat_ws(" ", col("first"), col("last"), col("maiden")).alias("full_name_safe"),
).show()
# +----------------+--------------+
# |full_name_unsafe|full_name_safe|
# +----------------+--------------+
# | Иван Петров| Иван Петров |
# | null| Анна Смирнова|
# +----------------+--------------+
concat_ws - стандартный инструмент для генерации бизнес-ключей и составных идентификаторов. Например, concat_ws("_", col("country"), col("region"), col("city")) → "RU_MSK_Moscow". Функция принимает как скалярные колонки, так и массивы (ArrayType) - разделитель вставляется между элементами массива.
lower, upper, initcap, trim: гигиена данных¶
Нормализация регистра и пробелов - обязательный шаг перед join-ами и дедупликацией. Два клиента "ИВАНОВ" и "иванов" без нормализации не совпадут при сравнении:
from pyspark.sql.functions import lower, upper, initcap, trim, ltrim, rtrim
df.select(
lower(col("name")).alias("lower"), # "иванов петров"
upper(col("name")).alias("upper"), # "ИВАНОВ ПЕТРОВ"
initcap(col("name")).alias("title"), # "Иванов Петров" (первая буква каждого слова)
trim(col("name")).alias("trimmed"), # убирает пробелы с обоих концов
ltrim(col("name")).alias("ltrimmed"), # только слева
rtrim(col("name")).alias("rtrimmed"), # только справа
).show()
initcap делает первую букву каждого слова заглавной - полезно для нормализации имён собственных, пришедших из разных источников. trim убирает не только пробелы, но и whitespace-символы (\t, \n). Spark 3.2+ поддерживает trim(col, chars) для удаления конкретных символов.
substring: нарезка строк на части¶
substring(col, pos, len) вырезает подстроку начиная с позиции pos длиной len символов. Позиции в Spark (как в SQL) начинаются с 1, а не с 0. Отрицательная позиция считает от конца строки:
from pyspark.sql.functions import substring
df = spark.createDataFrame([("7921234567",), ("74951234567",)], ["phone"])
df.select(
col("phone"),
substring(col("phone"), 2, 3).alias("operator_code"), # символы 2-4: "921", "495"
substring(col("phone"), -7, 7).alias("local_number"), # последние 7 цифр
).show()
# +-----------+-------------+------------+
# | phone|operator_code|local_number|
# +-----------+-------------+------------+
# | 7921234567| 921| 1234567|
# |74951234567| 495| 1234567|
# +-----------+-------------+------------+
substring - основной инструмент для парсинга данных с фиксированной шириной (fixed-width formats), распространённых в банковских и телекоммуникационных системах. left(col, n) и right(col, n) - удобные псевдонимы для substring с краёв строки (Spark 3.3+).
length и char_length: метрики строки¶
length() возвращает количество байт, char_length() - количество символов. Разница критична для Unicode: кириллица в UTF-8 занимает 2 байта на символ, поэтому для русского текста length вернёт вдвое большее число чем char_length:
from pyspark.sql.functions import length, char_length
df = spark.createDataFrame([("Иван",), ("Ivan",)], ["name"])
df.select(
col("name"),
length(col("name")).alias("bytes"), # Иван → 8, Ivan → 4
char_length(col("name")).alias("chars"), # Иван → 4, Ivan → 4
).show()
Используйте char_length для валидации данных: проверить что ИНН ровно 10 или 12 символов, что номер банковской карты 16 цифр, что СНИЛС 11 символов. length нужен для оценки объёма данных при хранении.
lpad и rpad: дополнение строк до фиксированной длины¶
lpad(col, len, pad) дополняет строку слева символами pad до длины len. rpad - то же самое справа. Если строка уже длиннее len - она обрезается. Незаменимы при работе с форматами фиксированной ширины: коды товаров, номера счетов, SWIFT-коды:
from pyspark.sql.functions import lpad, rpad
df = spark.createDataFrame([
("42",),
("1234",),
("987654",),
], ["code"])
df.select(
col("code"),
lpad(col("code"), 6, "0").alias("padded_left"), # "000042", "001234", "987654"
rpad(col("code"), 6, "X").alias("padded_right"), # "42XXXX", "1234XX", "987654"
).show()
# +------+-----------+------------+
# | code|padded_left|padded_right|
# +------+-----------+------------+
# | 42| 000042| 42XXXX|
# | 1234| 001234| 1234XX|
# |987654| 987654| 987654|
# +------+-----------+------------+
Типичный сценарий: нормализация кода продукта из legacy-системы, где поле имеет фиксированную ширину 8 символов с ведущими нулями. lpad(col("product_code"), 8, "0") приведёт "42" → "00000042" без Python UDF. Если строка длиннее целевой длины - lpad/rpad обрежет её: lpad("ABCDEF", 4, "0") → "ABCD".
Регулярные выражения: тяжёлая артиллерия¶
regexp_replace: очистка и стандартизация¶
regexp_replace(col, pattern, replacement) заменяет все вхождения паттерна на строку замены. Это основной инструмент нормализации: убрать лишние символы, привести к стандартному формату:
from pyspark.sql.functions import regexp_replace
df = spark.createDataFrame([
("+7 (921) 123-45-67",),
("8-495-123-45-67",),
("+7921 123 4567",),
], ["raw_phone"])
df.select(
col("raw_phone"),
# Шаг 1: убрать всё кроме цифр - класс [^\d] = "не цифра"
regexp_replace(col("raw_phone"), r"[^\d]", "").alias("digits_only"),
# Шаг 2: в цепочке - заменить ведущую 8 на 7 ($1 = первая группа)
regexp_replace(
regexp_replace(col("raw_phone"), r"[^\d]", ""),
r"^8(\d{10})$", "7$1"
).alias("normalized"),
).show(truncate=False)
# +-------------------+-----------+------------+
# |raw_phone |digits_only|normalized |
# +-------------------+-----------+------------+
# |+7 (921) 123-45-67 |79211234567|79211234567 |
# |8-495-123-45-67 |84951234567|74951234567 |
# |+7921 123 4567 |79211234567|79211234567 |
# +-------------------+-----------+------------+
Особенности Java-стиля regex в Spark (отличия от Python re):
- Backreferences в замене:
$1,$2(не\1,\2как в Python) \dработает, но в Python-строке нужно писать\\dили использовать raw stringr"\d"- Lookahead/lookbehind поддерживаются:
(?=...),(?<=...) - Флаги встраиваются в паттерн:
(?i)для case-insensitive,(?s)для dotall
regexp_extract: извлечение паттернов¶
regexp_extract(col, pattern, groupIndex) возвращает захватывающую группу с заданным индексом. Группа 0 - всё совпадение целиком. Если паттерн не совпадает - возвращает пустую строку (не NULL, в отличие от многих других функций):
from pyspark.sql.functions import regexp_extract
df = spark.createDataFrame([
("user=ivan@company.ru; action=login",),
("user=admin@corp.com; action=delete",),
("system: startup",), # нет email
], ["log_line"])
df.select(
col("log_line"),
# Группа 0 = всё совпадение (весь email)
regexp_extract(col("log_line"), r"[\w.+-]+@[\w-]+\.[a-z]+", 0).alias("email"),
# Группа 1 = первая захватывающая группа (только домен)
regexp_extract(col("log_line"), r"@([\w-]+\.[a-z]+)", 1).alias("domain"),
).show(truncate=False)
# +----------------------------------+---------------+-----------+
# |log_line |email |domain |
# +----------------------------------+---------------+-----------+
# |user=ivan@company.ru; action=login|ivan@company.ru|company.ru |
# |user=admin@corp.com; action=delete|admin@corp.com |corp.com |
# |system: startup | | |
# +----------------------------------+---------------+-----------+
Spark 3.1+ добавил regexp_extract_all - возвращает массив (ArrayType) всех совпадений, а не только первое. Незаменима когда в одной строке несколько значений одного паттерна:
from pyspark.sql.functions import regexp_extract_all
# Найти все числа в строке
df.select(
regexp_extract_all(col("text"), r"\d+", 0).alias("all_numbers")
)
# "order 123, item 456, qty 7" → ["123", "456", "7"]
Производительность регулярных выражений¶
Spark компилирует regex один раз на executor через java.util.regex.Pattern.compile и кэширует скомпилированный объект. Само сопоставление происходит на каждой строке.
Проблема - catastrophic backtracking: паттерны с .* и вложенными квантификаторами могут вызвать экспоненциальное время работы на аномальных строках. Например, (a+)+b на строке "aaaaaaaaaaaaaax" даёт тысячи итераций.
Практические правила:
- Используйте якоря
^и$когда нужна проверка всей строки - это устраняет backtracking на хвосте - Предпочитайте конкретные классы
[^\d]+вместо.*- они не допускают неоднозначности - Для проверки вхождения подстроки - используйте
contains(), не regex - Тестируйте regex на строках-аномалиях (очень длинных или содержащих повторяющиеся символы)
split и работа с массивами¶
split(col, pattern, limit=-1) делит строку по разделителю и возвращает ArrayType(StringType). Это ключевой шаг при парсинге delimiter-based форматов: CSV-в-строке, pipe-separated, tag-списки:
from pyspark.sql.functions import split, element_at, size, array_join
df = spark.createDataFrame([
("RU|Moscow|Central",),
("DE|Berlin|East",),
("US|New York|Northeast",),
], ["location"])
result = df.select(
col("location"),
split(col("location"), r"\|").alias("parts"), # Array["RU","Moscow","Central"]
# element_at использует 1-based индексы (как SQL)
element_at(split(col("location"), r"\|"), 1).alias("country"), # "RU"
element_at(split(col("location"), r"\|"), 2).alias("city"), # "Moscow"
element_at(split(col("location"), r"\|"), -1).alias("region"), # "Central" (от конца)
size(split(col("location"), r"\|")).alias("part_count"), # 3
)
result.show(truncate=False)
element_at использует 1-based индексы (как SQL), отрицательный индекс считает с конца. Альтернативный синтаксис: split(col, "|")[0] - возвращает элемент в Python-style (0-based), возвращает NULL при выходе за границы вместо исключения как у element_at.
array_join - обратная операция: собрать массив обратно в строку с разделителем:
from pyspark.sql.functions import array_join
# Array["RU","Moscow"] → "RU/Moscow"
df.select(array_join(split(col("location"), r"\|"), "/").alias("reformatted"))
# Третий аргумент - замена для NULL в массиве:
# array_join(arr, ",", "N/A") → null-элементы заменяются на "N/A"
array_join полезна когда нужно реформатировать разделитель или собрать теги/категории из массива обратно в строку.
Поиск и предикаты¶
contains, startswith, endswith¶
Методы Column-объекта для проверки вхождения без regex-оверхеда. Транслируются в простые JVM-операции String.startsWith(), String.contains():
# Фильтрация - Column methods
df.filter(col("url").startswith("https://"))
df.filter(col("email").endswith(".ru"))
df.filter(col("user_agent").contains("Mozilla"))
# В select - возвращают BooleanType
df.select(
col("url").startswith("https").alias("is_secure"),
col("domain").endswith(".gov").alias("is_government"),
)
Для простых проверок на вхождение подстроки используйте эти методы вместо rlike - они быстрее и читабельнее.
instr и locate: позиция подстроки¶
instr(col, substr) возвращает позицию первого вхождения (1-based, 0 если не найдено). Полезно когда нужна не просто проверка, а позиция для последующей обработки:
from pyspark.sql.functions import instr
df.select(
col("text"),
instr(col("text"), "@").alias("at_pos"), # позиция символа @
(instr(col("text"), "@") > 0).alias("has_at"), # булевая проверка
# Извлечь часть строки до символа используя instr
substring(col("text"), 1, instr(col("text"), "@") - 1).alias("before_at"),
)
locate(substr, col, pos=1) - то же самое, но можно задать начальную позицию поиска (для нахождения второго вхождения).
levenshtein: расстояние редактирования¶
levenshtein(col1, col2) вычисляет расстояние Левенштейна - минимальное количество односимвольных операций (вставка, удаление, замена), необходимых для превращения одной строки в другую. Используется для нечёткого сравнения строк (fuzzy matching) при дедупликации и поиске похожих записей:
from pyspark.sql.functions import levenshtein, col, lit
df = spark.createDataFrame([
("Иванов", "Иванов"), # точное совпадение
("Иванов", "Ивановв"), # опечатка (лишняя буква)
("Иванов", "Иванова"), # одна замена
("Иванов", "Петров"), # разные фамилии
], ["name_a", "name_b"])
df.select(
col("name_a"),
col("name_b"),
levenshtein(col("name_a"), col("name_b")).alias("distance"),
).show()
# +-------+--------+--------+
# | name_a| name_b|distance|
# +-------+--------+--------+
# |Иванов | Иванов| 0|
# |Иванов | Ивановв| 1|
# |Иванов | Иванова| 1|
# |Иванов | Петров| 5|
# +-------+--------+--------+
# Фильтрация потенциальных дублей: distance <= 2
df.filter(levenshtein(col("name_a"), col("name_b")) <= 2).show()
В production используется для дедупликации клиентских баз: найти записи с небольшими опечатками в имени или email. Важно об оптимизации: levenshtein вычисляется попарно для каждой строки - он применяется к уже объединённым через join записям, а не как полный cross join всего датасета. Порог <= 2 покрывает большинство опечаток при вводе, порог <= 3 начинает давать ложные срабатывания на коротких именах.
like и rlike¶
like использует SQL LIKE-паттерны: % - любая строка, _ - один символ. Быстрее регулярок для простых wildcard-проверок. rlike принимает полноценное регулярное выражение:
df.filter(col("name").like("Ива%")) # начинается с "Ива"
df.filter(col("code").like("RU__%")) # "RU" + минимум 2 символа
df.filter(col("phone").rlike(r"^\+7\d{10}$")) # полный regex
NULL propagation в строковых функциях¶
Большинство строковых функций следуют принципу NULL propagation: если хотя бы один аргумент NULL - результат NULL. Это стандарт SQL ANSI, но легко приводит к неожиданным потерям данных в production:
Один NULL-элемент полностью «заражает» результат concat. concat_ws намеренно спроектирована иначе: она пропускает NULL-аргументы и не добавляет лишние разделители. Для других функций используйте coalesce(col("name"), lit("")) перед конкатенацией.
Функции length(), lower(), upper(), trim() возвращают NULL на NULL-входе. regexp_extract при несовпадении паттерна возвращает "" (пустую строку), а не NULL - это важно при последующей фильтрации.
replace vs regexp_replace¶
replace(col, search, replacement) (Spark 3.4+) - буквальная замена строки без компиляции регулярного выражения. Существенно быстрее regexp_replace для статических паттернов:
from pyspark.sql.functions import replace
df.select(
# Простая замена - используй replace
replace(col("url"), lit("http://"), lit("https://")).alias("secured"),
# С паттерном - только тогда regexp_replace
regexp_replace(col("phone"), r"[^\d]", "").alias("digits_only"),
)
replace транслируется в String.replace() без regex overhead. Используйте его для замены конкретных подстрок и regexp_replace только когда нужна мощь паттернов.
Higher-order functions как граница между built-in и UDF¶
Когда встроенных функций не хватает для обработки массивов, рассмотрите higher-order functions как промежуточный вариант перед UDF. Они выполняются в JVM, а Catalyst анализирует лямбда-выражение:
from pyspark.sql.functions import transform, filter as array_filter, lower, trim, initcap
df = spark.createDataFrame([
(["Иванов", "петров", " СИДОРОВ "],),
(["Кузнецов", None, "Новиков"],),
], ["names"])
df.select(
# Нормализация каждого элемента массива без UDF
transform(
col("names"),
lambda x: initcap(trim(lower(x)))
).alias("normalized_names"),
# Фильтрация NULL из массива
array_filter(
col("names"),
lambda x: x.isNotNull()
).alias("non_null_names"),
).show(truncate=False)
# +-------------------------+------------------+
# |normalized_names |non_null_names |
# +-------------------------+------------------+
# |[Иванов, Петров, Сидоров]|[Иванов, петров,.]|
# |[Кузнецов, null, Новиков]|[Кузнецов, Новиков]|
# +-------------------------+------------------+
transform применяет функцию к каждому элементу массива. Лямбда получает Column-объект, значит внутри можно использовать любые встроенные функции. Это работает через JVM и значительно быстрее Python UDF, хотя всё равно медленнее специализированных функций типа array_distinct.
Explain plan для строковых трансформаций¶
Посмотрим как Spark компилирует цепочку встроенных функций:
df.select(
regexp_replace(
trim(lower(col("phone"))),
r"[^\d]", ""
).alias("clean_phone")
).explain()
== Physical Plan ==
*(1) Project [regexp_replace(trim(lower(phone#0)), [^\d], , 1) AS clean_phone#3]
+- *(1) Scan ExistingRDD[phone#0]
Ключевое наблюдение: весь pipeline в одном *(1) - WholeStageCodeGen. Spark объединяет lower → trim → regexp_replace в единый JVM-байткод проход без промежуточных объектов в heap.
Сравните с Python UDF:
from pyspark.sql.functions import udf
from pyspark.sql.types import StringType
import re
@udf(StringType())
def clean_phone_udf(s):
return re.sub(r"[^\d]", "", s.lower().strip()) if s else None
df.select(clean_phone_udf(col("phone"))).explain()
== Physical Plan ==
*(2) Project [pythonUDF0#5 AS clean_phone_udf(phone)#6]
+- BatchEvalPython [clean_phone_udf(phone#0)], [pythonUDF0#5] ← Python barrier
+- *(1) Scan ExistingRDD[phone#0]
Появляется BatchEvalPython - узел, который сериализует батч строк, отправляет в Python-процесс и ждёт результата. WholeStageCodeGen для этого шага выключается. Дополнительно: Catalyst не может применить к нему constant folding или predicate pushdown.
Bronze → Silver: нормализация клиентских данных¶
Полный pipeline нормализации сырых клиентских данных - типичная задача при переходе из Bronze в Silver layer:
from pyspark.sql.functions import (
col, lower, trim, regexp_replace, regexp_extract,
concat_ws, when, lit, char_length
)
raw_customers = spark.read.parquet("s3://lake/bronze/customers/")
silver_customers = raw_customers.select(
col("customer_id"),
# Нормализация имени: убрать пробелы, привести к lower
trim(lower(col("first_name"))).alias("first_name"),
trim(lower(col("last_name"))).alias("last_name"),
# Телефон: только цифры, заменить ведущую 8 → 7
# Шаг 1: убрать все не-цифры
# Шаг 2: заменить ведущую 8 на 7 для 11-значных номеров
regexp_replace(
regexp_replace(col("phone"), r"[^\d]", ""),
r"^8(\d{10})$", "7$1"
).alias("phone_normalized"),
# Email: lower + trim, проверить формат через regexp_extract
when(
regexp_extract(
trim(lower(col("email"))),
r"^[\w.+-]+@[\w-]+\.[a-z]{2,}$", 0
) != "",
trim(lower(col("email")))
).otherwise(lit(None)).alias("email_normalized"),
# ИНН: валидация длины (10 или 12 символов, только цифры)
when(
char_length(col("inn")).isin(10, 12) &
(regexp_extract(col("inn"), r"^\d+$", 0) != ""),
col("inn")
).otherwise(lit(None)).alias("inn_validated"),
# Бизнес-ключ для дедупликации - concat_ws автоматически пропускает NULL
concat_ws(
"_",
trim(lower(col("first_name"))),
trim(lower(col("last_name"))),
regexp_replace(col("phone"), r"[^\d]", "")
).alias("dedup_key"),
col("created_at"),
)
silver_customers.write.mode("overwrite").parquet("s3://lake/silver/customers/")
Весь pipeline выполняется в JVM без Python-сериализации. Catalyst объединяет все трансформации в минимальное количество проходов по данным - в физическом плане будет один *(1) Project узел.
Практика¶
Задание 1: Нормализация email-адресов¶
Дан DataFrame с колонкой raw_email. Нужно:
- Привести к нижнему регистру и убрать пробелы
- Проверить формат через
regexp_extract - Извлечь домен в отдельную колонку
- Записать
NULLдля некорректных email
df = spark.createDataFrame([
(" Ivan@COMPANY.RU ",),
("admin@corp.com",),
("not-an-email",),
(None,),
], ["raw_email"])
Решение
from pyspark.sql.functions import trim, lower, regexp_extract, when, lit, col
EMAIL_RE = r"^([\w.+-]+)@([\w-]+\.[a-z]{2,})$"
result = df.select(
col("raw_email"),
when(
regexp_extract(trim(lower(col("raw_email"))), EMAIL_RE, 0) != "",
trim(lower(col("raw_email")))
).otherwise(lit(None)).alias("email"),
regexp_extract(trim(lower(col("raw_email"))), EMAIL_RE, 2).alias("domain"),
)
result.show(truncate=False)
# +-----------------+---------------+----------+
# |raw_email |email |domain |
# +-----------------+---------------+----------+
# | Ivan@COMPANY.RU|ivan@company.ru|company.ru|
# |admin@corp.com |admin@corp.com |corp.com |
# |not-an-email |null | |
# |null |null | |
# +-----------------+---------------+----------+
regexp_extract возвращает "" (не NULL) при несовпадении, поэтому проверяем != "". Для NULL в raw_email функции вернут NULL, и when тоже вернёт NULL через otherwise(lit(None)).
Задание 2: Парсинг лога в Data Lake¶
Входная строка лога имеет формат:
USER_123|2023-01-15|192.168.1.1|Mozilla/5.0 (Windows NT 10.0)
Нужно распарсить в структурированные поля: user_id (Integer), event_date (String), ip_address, browser_family (только первое слово до /).
df = spark.createDataFrame([
("USER_123|2023-01-15|192.168.1.1|Mozilla/5.0 (Windows NT 10.0)",),
("USER_456|2023-01-15|10.0.0.1|curl/7.68.0",),
("USER_789|2023-01-16|172.16.0.5|Python-requests/2.28",),
], ["raw_log"])
Решение
from pyspark.sql.functions import split, element_at, regexp_extract, col
from pyspark.sql.types import IntegerType
# Вычисляем split один раз и переиспользуем через withColumn
parsed = df.withColumn("parts", split(col("raw_log"), r"\|"))
result = parsed.select(
# user_id: убрать "USER_" и кастовать в int
regexp_extract(element_at(col("parts"), 1), r"USER_(\d+)", 1)
.cast(IntegerType()).alias("user_id"),
# date: второй элемент напрямую
element_at(col("parts"), 2).alias("event_date"),
# IP: третий элемент
element_at(col("parts"), 3).alias("ip_address"),
# browser: от начала строки до символа "/" (исключая версию)
regexp_extract(element_at(col("parts"), 4), r"^([^/]+)", 1).alias("browser_family"),
)
result.show(truncate=False)
# +-------+----------+-----------+----------------+
# |user_id|event_date|ip_address |browser_family |
# +-------+----------+-----------+----------------+
# | 123|2023-01-15|192.168.1.1|Mozilla/5.0 ... |
# | 456|2023-01-15|10.0.0.1 |curl |
# | 789|2023-01-16|172.16.0.5 |Python-requests |
# +-------+----------+-----------+----------------+
split вызывается один раз через withColumn("parts", ...), затем переиспользуется. Если бы мы писали split(col("raw_log"), r"\|") в каждом element_at - Spark выполнял бы split 4 раза. Catalyst в текущих версиях не всегда объединяет одинаковые выражения автоматически.
Задание 3: Генерация бизнес-ключа для дедупликации¶
Дан DataFrame с клиентскими данными. Нужно создать dedup_key из нормализованных полей. NULL-поля должны исключаться из ключа. Результат должен быть однозначным независимо от пробелов в исходных данных:
df = spark.createDataFrame([
("Иван", "Петров", "+7 (921) 123-45-67", "ivan@mail.ru"),
("Анна", "Петрова", None, "anna@mail.ru"),
(" иван ", " Петров ", "7-921-123-45-67", None),
], ["first_name", "last_name", "phone", "email"])
Решение
from pyspark.sql.functions import concat_ws, lower, trim, regexp_replace, col
result = df.select(
concat_ws(
"|",
trim(lower(col("first_name"))),
trim(lower(col("last_name"))),
regexp_replace(col("phone"), r"[^\d]", ""), # NULL → NULL (пропускается concat_ws)
trim(lower(col("email"))),
).alias("dedup_key")
)
result.show(truncate=False)
# +------------------------------------------+
# |dedup_key |
# +------------------------------------------+
# |иван|петров|79211234567|ivan@mail.ru |
# |анна|петрова|anna@mail.ru |
# |иван|петров|79211234567 |
# +------------------------------------------+
concat_ws автоматически пропускает NULL, поэтому строки 2 и 3 формируют корректные ключи без лишних разделителей. Строки 1 и 3 получат одинаковый ключ по телефону, несмотря на разные форматы и пробелы в именах - это позволяет найти дублирующиеся записи.
Задание 4: Валидация качества данных¶
Для DataFrame с колонкой inn (ИНН) и phone нужно:
- Проверить ИНН: 10 или 12 символов, только цифры
- Проверить телефон: после нормализации должно быть ровно 11 цифр, начинаться с 7
- Создать поле
quality_flagсо значениямиvalid,partial(один из двух некорректен),invalid
Решение
from pyspark.sql.functions import (
col, char_length, regexp_extract, regexp_replace, when
)
df_quality = df.select(
col("customer_id"),
col("inn"),
col("phone"),
# Нормализуем телефон для проверки
regexp_replace(col("phone"), r"[^\d]", "").alias("phone_digits"),
).select(
col("customer_id"),
# Флаг валидности ИНН
(
char_length(col("inn")).isin(10, 12) &
(regexp_extract(col("inn"), r"^\d+$", 0) != "")
).alias("inn_ok"),
# Флаг валидности телефона
(
(char_length(col("phone_digits")) == 11) &
col("phone_digits").startswith("7")
).alias("phone_ok"),
).select(
col("customer_id"),
when(col("inn_ok") & col("phone_ok"), "valid")
.when(col("inn_ok") | col("phone_ok"), "partial")
.otherwise("invalid").alias("quality_flag")
)
df_quality.groupBy("quality_flag").count().show()
Сравнение производительности¶
На 100 млн строк при нормализации телефонного номера:
| Метод | Время | Примечание |
|---|---|---|
regexp_replace (built-in) |
~8 сек | WholeStageCodeGen, JVM bytecode |
| Pandas UDF (Arrow-based) | ~45 сек | Батчи Arrow, пересекает JVM→Python |
| Python UDF (обычный) | ~4 мин | Row-by-row pickle, нет CodeGen |
Разница в 30× между built-in и обычным UDF - не исключение, а закономерность для строковых операций. В production pipeline с нормализацией десятков полей это может означать разницу между 5 минутами и 2 часами работы джобы.
Итоговое правило: если задача решается через pyspark.sql.functions - используй их. Если нужна обработка массивов - рассмотри transform/filter/aggregate (higher-order functions). Python UDF - только когда нет встроенной альтернативы (например, специфичная кодировка, криптография, внешняя библиотека).