Строковые функции в PySpark

concat, concat_ws, substring, regexp_extract, instr, regexp_replace, upper/lower/initcap, trim/ltrim/rtrim, lpad/rpad, length, levenshtein.

core

Строковые функции PySpark позволяют обрабатывать и преобразовывать текстовые данные. Они особенно полезны при очистке данных, извлечении информации и трансформации текстовых колонок. Большинство из них знакомы SQL-разработчикам — PySpark предоставляет к ним доступ через модуль pyspark.sql.functions.


Основные строковые функции

1. Конкатенация

  • concat(*cols): Объединяет несколько колонок или строк в одну.
  • concat_ws(delimiter, *cols): Объединяет колонки или строки с указанным разделителем.
from pyspark.sql.functions import concat, concat_ws

# Объединить имя и фамилию
df = df.withColumn("full_name", concat(df.first_name, df.last_name))

# Объединить с пробелом в качестве разделителя
df = df.withColumn("full_name_with_space", concat_ws(" ", df.first_name, df.last_name))

2. Подстрока и извлечение

  • substring(col, pos, length): Извлекает подстроку из колонки.
  • substr(col, pos, length): Псевдоним для substring.
  • regexp_extract(col, pattern, groupIdx): Извлекает совпадение из строки по regex-паттерну.
from pyspark.sql.functions import substring, regexp_extract

# Первые 5 символов полного имени
df = df.withColumn("name_substring", substring(df.full_name, 1, 5))

# Домен из email через regex
df = df.withColumn("email_domain", regexp_extract(df.email, r"@(\w+)", 1))

3. Поиск и замена

  • instr(col, substring): Находит позицию первого вхождения подстроки.
  • regexp_replace(col, pattern, replacement): Заменяет подстроки, соответствующие regex-паттерну.
from pyspark.sql.functions import instr, regexp_replace

# Позиция первого вхождения 'error'
df = df.withColumn("error_position", instr(df.description, "error"))

# Удалить все цифры из текстовой колонки
df = df.withColumn("text_without_digits", regexp_replace(df.text, r"[0-9]", ""))

4. Изменение регистра

  • upper(col): Переводит текст в верхний регистр.
  • lower(col): Переводит текст в нижний регистр.
  • initcap(col): Делает первую букву каждого слова заглавной.
from pyspark.sql.functions import upper, lower, initcap

# Имя в верхнем регистре
df = df.withColumn("first_name_upper", upper(df.first_name))

# Имя в нижнем регистре
df = df.withColumn("first_name_lower", lower(df.first_name))

# Первая буква каждого слова заглавная
df = df.withColumn("name_initcap", initcap(df.full_name))

5. Обрезка и дополнение

  • trim(col): Удаляет пробелы в начале и конце строки.
  • ltrim(col): Удаляет пробелы в начале строки.
  • rtrim(col): Удаляет пробелы в конце строки.
  • lpad(col, length, pad): Дополняет строку символами слева до заданной длины.
  • rpad(col, length, pad): Дополняет строку символами справа до заданной длины.
from pyspark.sql.functions import trim, lpad, rpad

# Убрать пробелы в начале и конце имени пользователя
df = df.withColumn("trimmed_username", trim(df.username))

# Дополнить ID аккаунта нулями слева до длины 10
df = df.withColumn("padded_account_id", lpad(df.account_id, 10, "0"))

# Дополнить код символами "#" справа до длины 8
df = df.withColumn("padded_code", rpad(df.code, 8, "#"))

6. Длина и сравнение строк

  • length(col): Возвращает длину строки.
  • levenshtein(col1, col2): Вычисляет расстояние Левенштейна (редакционное расстояние) между двумя строками.
from pyspark.sql.functions import length, levenshtein

# Длина пароля
df = df.withColumn("password_length", length(df.password))

# Расстояние Левенштейна между двумя строками
df = df.withColumn("levenshtein_distance", levenshtein(df.string1, df.string2))

Полный пример

from pyspark.sql import SparkSession
from pyspark.sql.functions import (
    col, concat, concat_ws, substring, upper, lower, initcap,
    trim, ltrim, rtrim, regexp_replace, regexp_extract, length,
    instr, lpad, rpad
)

# Тестовые данные
data = [
    ("John", "Doe", "john.doe@example.com", "   active   ", "12345"),
    ("Jane", "Smith", "jane.smith@work.org", "inactive", "67890"),
    ("Sam", "Brown", "sam.brown@data.net", "active   ", "111213")
]

# Создание DataFrame
spark = SparkSession.builder.appName("StringFunctionsWithColumn").getOrCreate()
df = spark.createDataFrame(data, ["first_name", "last_name", "email", "status", "account_id"])

# Применение строковых функций через withColumn
df_transformed = df \
    .withColumn("full_name", concat_ws(" ", col("first_name"), col("last_name"))) \
    .withColumn("email_uppercase", upper(col("email"))) \
    .withColumn("email_domain", regexp_extract(col("email"), r"@(\w+)", 1)) \
    .withColumn("trimmed_status", trim(col("status"))) \
    .withColumn("padded_account_id", lpad(col("account_id"), 10, "0")) \
    .withColumn("email_prefix", substring(col("email"), 1, 5)) \
    .withColumn("cleaned_email", regexp_replace(col("email"), r"[.@]", "-")) \
    .withColumn("email_length", length(col("email"))) \
    .withColumn("first_name_initcap", initcap(col("first_name"))) \
    .withColumn("description", concat_ws(" | ", col("full_name"), col("trimmed_status"), col("email_domain")))

# Показать результирующий DataFrame
df_transformed.show(truncate=False)