Строковые функции в 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)