String Functions: Spark SQL против Python UDF

Встроенные строковые функции PySpark - regexp_replace, concat_ws, split, substring и их преимущества над Python UDF в production ETL.

core

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:

Что происходит для каждой строки датасета:

  1. JVM сериализует UnsafeRow в Python-объект через Pickle - дорогостоящая операция, требующая копирования данных
  2. Данные отправляются в отдельный Python-процесс (он запускается рядом с JVM-воркером)
  3. Python выполняет вашу функцию и формирует результат
  4. Результат сериализуется обратно в байтовый поток (Pickle)
  5. 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 string r"\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. Нужно:

  1. Привести к нижнему регистру и убрать пробелы
  2. Проверить формат через regexp_extract
  3. Извлечь домен в отдельную колонку
  4. Записать 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 нужно:

  1. Проверить ИНН: 10 или 12 символов, только цифры
  2. Проверить телефон: после нормализации должно быть ровно 11 цифр, начинаться с 7
  3. Создать поле 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 - только когда нет встроенной альтернативы (например, специфичная кодировка, криптография, внешняя библиотека).