101 практическое упражнение по PySpark

68 задач с решениями: создание DataFrame, фильтрация, агрегации, оконные функции, ML-пайплайны, работа с датами, UDF, pivot/unpivot, нормализация и многое другое.

core practice

Этот урок содержит набор практических задач по PySpark для закрепления навыков анализа данных. Упражнения разбиты по уровням сложности: L1 — начальный, L2 — средний, L3 — продвинутый. Каждая задача снабжена входными данными и решением. Рекомендуется сначала попытаться решить задачу самостоятельно, а затем сверяться с ответом.

Инициализация SparkSession (выполните перед всеми упражнениями):

import findspark
findspark.init()

from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("PySpark 101 Exercises").getOrCreate()

1. Как импортировать PySpark и проверить версию?

Уровень: L1

import findspark
findspark.init()

# SparkSession — точка входа для работы с DataFrame и SQL API в PySpark.
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("PySpark 101 Exercises").getOrCreate()

# Проверяем версию
print(spark.version)
#> 3.3.2

2. Как преобразовать индекс PySpark DataFrame в столбец?

Уровень: L1

Подсказка: В PySpark DataFrame нет явного понятия индекса, как в Pandas. Однако можно добавить новый столбец с номером строки.

Входные данные:

df = spark.createDataFrame([
    ("Alice", 1),
    ("Bob", 2),
    ("Charlie", 3),
], ["Name", "Value"])

df.show()
+-------+-----+
|   Name|Value|
+-------+-----+
|  Alice|    1|
|    Bob|    2|
|Charlie|    3|
+-------+-----+

Ожидаемый результат:

+-------+-----+-----+
|   Name|Value|index|
+-------+-----+-----+
|  Alice|    1|    0|
|    Bob|    2|    1|
|Charlie|    3|    2|
+-------+-----+-----+
Решение
from pyspark.sql.window import Window
from pyspark.sql.functions import row_number, monotonically_increasing_id

# Определяем оконную спецификацию
w = Window.orderBy(monotonically_increasing_id())

# Добавляем индекс (нумерация с 0)
df = df.withColumn("index", row_number().over(w) - 1)

df.show()
+-------+-----+-----+
|   Name|Value|index|
+-------+-----+-----+
|  Alice|    1|    0|
|    Bob|    2|    1|
|Charlie|    3|    2|
+-------+-----+-----+

3. Как создать PySpark DataFrame из нескольких списков?

Уровень: L1

Создайте PySpark DataFrame из list1 и list2.

Подсказка: Сначала создайте RDD из списков, затем преобразуйте его в DataFrame.

Входные данные:

list1 = ["a", "b", "c", "d"]
list2 = [1, 2, 3, 4]
Решение
# Создаём RDD из списков и преобразуем в DataFrame
rdd = spark.sparkContext.parallelize(list(zip(list1, list2)))
df = rdd.toDF(["Column1", "Column2"])

df.show()
+-------+-------+
|Column1|Column2|
+-------+-------+
|      a|      1|
|      b|      2|
|      c|      3|
|      d|      4|
+-------+-------+

4. Как получить элементы списка A, отсутствующие в списке B?

Уровень: L2

Получите элементы list_A, которых нет в list_B. Используйте операцию subtract на RDD.

Входные данные:

list_A = [1, 2, 3, 4, 5]
list_B = [4, 5, 6, 7, 8]

Ожидаемый результат: [1, 2, 3]

Решение
sc = spark.sparkContext

# Преобразуем списки в RDD
rdd_A = sc.parallelize(list_A)
rdd_B = sc.parallelize(list_B)

# Операция subtract
result_rdd = rdd_A.subtract(rdd_B)

result_list = result_rdd.collect()
print(result_list)
#> [1, 2, 3]

5. Как получить элементы, не общие для обоих списков A и B?

Уровень: L2

Получите все элементы list_A и list_B, которые не входят в их пересечение (симметрическая разность).

Входные данные:

list_A = [1, 2, 3, 4, 5]
list_B = [4, 5, 6, 7, 8]
Решение
sc = spark.sparkContext

rdd_A = sc.parallelize(list_A)
rdd_B = sc.parallelize(list_B)

# Элементы A, которых нет в B, и элементы B, которых нет в A
result_rdd_A = rdd_A.subtract(rdd_B)
result_rdd_B = rdd_B.subtract(rdd_A)

# Объединяем два RDD
result_rdd = result_rdd_A.union(result_rdd_B)

result_list = result_rdd.collect()
print(result_list)
[1, 2, 3, 8, 6, 7]

6. Как получить минимум, 25-й процентиль, медиану, 75-й и максимум числового столбца?

Уровень: L2

Вычислите минимум, 25-й процентиль, медиану, 75-й процентиль и максимум столбца Age.

Входные данные:

data = [("A", 10), ("B", 20), ("C", 30), ("D", 40), ("E", 50),
        ("F", 15), ("G", 28), ("H", 54), ("I", 41), ("J", 86)]
df = spark.createDataFrame(data, ["Name", "Age"])

df.show()
+----+---+
|Name|Age|
+----+---+
|   A| 10|
|   B| 20|
|   C| 30|
|   D| 40|
|   E| 50|
|   F| 15|
|   G| 28|
|   H| 54|
|   I| 41|
|   J| 86|
+----+---+
Решение
# Вычисляем процентили
quantiles = df.approxQuantile("Age", [0.0, 0.25, 0.5, 0.75, 1.0], 0.01)

print("Min: ",              quantiles[0])
print("25-й процентиль: ", quantiles[1])
print("Медиана: ",         quantiles[2])
print("75-й процентиль: ", quantiles[3])
print("Max: ",             quantiles[4])
Min:  10.0
25-й процентиль:  20.0
Медиана:  30.0
75-й процентиль:  50.0
Max:  86.0

7. Как получить частоту уникальных значений в столбце?

Уровень: L1

Подсчитайте частоту встречаемости каждого уникального значения.

Входные данные:

from pyspark.sql import Row

data = [
    Row(name='John', job='Engineer'),
    Row(name='John', job='Engineer'),
    Row(name='Mary', job='Scientist'),
    Row(name='Bob',  job='Engineer'),
    Row(name='Bob',  job='Engineer'),
    Row(name='Bob',  job='Scientist'),
    Row(name='Sam',  job='Doctor'),
]

df = spark.createDataFrame(data)
df.show()
+----+---------+
|name|      job|
+----+---------+
|John| Engineer|
|John| Engineer|
|Mary|Scientist|
| Bob| Engineer|
| Bob| Engineer|
| Bob|Scientist|
| Sam|   Doctor|
+----+---------+
Решение
df.groupBy("job").count().show()
+---------+-----+
|      job|count|
+---------+-----+
| Engineer|    4|
|Scientist|    2|
|   Doctor|    1|
+---------+-----+

8. Как сохранить только 2 наиболее частых значения и заменить остальные на 'Other'?

Уровень: L3

Входные данные:

from pyspark.sql import Row

data = [
    Row(name='John', job='Engineer'),
    Row(name='John', job='Engineer'),
    Row(name='Mary', job='Scientist'),
    Row(name='Bob',  job='Engineer'),
    Row(name='Bob',  job='Engineer'),
    Row(name='Bob',  job='Scientist'),
    Row(name='Sam',  job='Doctor'),
]

df = spark.createDataFrame(data)
df.show()
+----+---------+
|name|      job|
+----+---------+
|John| Engineer|
|John| Engineer|
|Mary|Scientist|
| Bob| Engineer|
| Bob| Engineer|
| Bob|Scientist|
| Sam|   Doctor|
+----+---------+
Решение
from pyspark.sql.functions import col, when

# Получаем 2 наиболее частых значения
top_2_jobs = (df.groupBy('job')
               .count()
               .orderBy('count', ascending=False)
               .limit(2)
               .select('job')
               .rdd.flatMap(lambda x: x)
               .collect())

# Заменяем все остальные значения на 'Other'
df = df.withColumn('job',
    when(col('job').isin(top_2_jobs), col('job')).otherwise('Other'))

df.show()
+----+---------+
|name|      job|
+----+---------+
|John| Engineer|
|John| Engineer|
|Mary|Scientist|
| Bob| Engineer|
| Bob| Engineer|
| Bob|Scientist|
| Sam|    Other|
+----+---------+

9. Как удалить строки с NA-значениями в конкретном столбце?

Уровень: L1

Входные данные:

df = spark.createDataFrame([
    ("A",  1,    None),
    ("B",  None, "123"),
    ("B",  3,    "456"),
    ("D",  None, None),
], ["Name", "Value", "id"])

df.show()
+----+-----+----+
|Name|Value|  id|
+----+-----+----+
|   A|    1|null|
|   B| null| 123|
|   B|    3| 456|
|   D| null|null|
+----+-----+----+
Решение
df_2 = df.dropna(subset=['Value'])
df_2.show()
+----+-----+----+
|Name|Value|  id|
+----+-----+----+
|   A|    1|null|
|   B|    3| 456|
+----+-----+----+

10. Как переименовать столбцы DataFrame с помощью двух списков?

Уровень: L2

Переименуйте столбцы, используя список старых имён и список новых имён.

Входные данные:

df = spark.createDataFrame([(1, 2, 3), (4, 5, 6)], ["col1", "col2", "col3"])

old_names = ["col1", "col2", "col3"]
new_names = ["new_col1", "new_col2", "new_col3"]

df.show()
+----+----+----+
|col1|col2|col3|
+----+----+----+
|   1|   2|   3|
|   4|   5|   6|
+----+----+----+
Решение
for old_name, new_name in zip(old_names, new_names):
    df = df.withColumnRenamed(old_name, new_name)

df.show()
+--------+--------+--------+
|new_col1|new_col2|new_col3|
+--------+--------+--------+
|       1|       2|       3|
|       4|       5|       6|
+--------+--------+--------+

11. Как разбить числовой список на 10 групп равного размера?

Уровень: L2

Входные данные:

from pyspark.sql.functions import rand
from pyspark.ml.feature import Bucketizer

num_items = 100
df = spark.range(num_items).select(rand(seed=42).alias("values"))

df.show(5)
+-------------------+
|             values|
+-------------------+
| 0.619189370225301|
|0.5096018842446481|
|0.8325259388871524|
|0.26322809041172357|
|0.6702867696264135|
+-------------------+
only showing top 5 rows
Решение
# Определяем границы бакетов
num_buckets = 10
quantiles = df.stat.approxQuantile("values",
    [i/num_buckets for i in range(num_buckets + 1)], 0.01)

# Создаём Bucketizer
bucketizer = Bucketizer(splits=quantiles, inputCol="values", outputCol="buckets")

# Применяем
df_buck = bucketizer.transform(df)

# Таблица частот
df_buck.groupBy("buckets").count().show()

df_buck.show(5)
+-------+-----+
|buckets|count|
+-------+-----+
|    8.0|   10|
|    0.0|    8|
|    7.0|   10|
|    1.0|   10|
|    4.0|   10|
|    3.0|   10|
|    2.0|   10|
|    6.0|   10|
|    5.0|   10|
|    9.0|   12|
+-------+-----+

12. Как создать таблицу сопряжённости (contingency table)?

Уровень: L1

Входные данные:

data = [("A", "X"), ("A", "Y"), ("A", "X"),
        ("B", "Y"), ("B", "X"),
        ("C", "X"), ("C", "X"), ("C", "Y")]
df = spark.createDataFrame(data, ["category1", "category2"])

df.show()
+---------+---------+
|category1|category2|
+---------+---------+
|        A|        X|
|        A|        Y|
|        A|        X|
|        B|        Y|
|        B|        X|
|        C|        X|
|        C|        X|
|        C|        Y|
+---------+---------+
Решение
# Частоты
df.cube("category1").count().show()

# Таблица сопряжённости
df.crosstab('category1', 'category2').show()
+---------+-----+
|category1|count|
+---------+-----+
|     null|    8|
|        A|    3|
|        B|    2|
|        C|    3|
+---------+-----+

+-------------------+---+---+
|category1_category2|  X|  Y|
+-------------------+---+---+
|                  A|  2|  1|
|                  B|  1|  1|
|                  C|  2|  1|
+-------------------+---+---+

13. Как найти числа, кратные 3, в столбце?

Уровень: L2

Входные данные:

from pyspark.sql.functions import rand

df = spark.range(10)
df = df.withColumn("random", ((rand(seed=42) * 10) + 1).cast("int"))

df.show()
+---+------+
| id|random|
+---+------+
|  0|     7|
|  1|     6|
|  2|     9|
|  3|     7|
|  4|     3|
|  5|     8|
|  6|     9|
|  7|     8|
|  8|     3|
|  9|     8|
+---+------+
Решение
from pyspark.sql.functions import col, when

df = df.withColumn("is_multiple_of_3",
    when(col("random") % 3 == 0, 1).otherwise(0))

df.show()
+---+------+----------------+
| id|random|is_multiple_of_3|
+---+------+----------------+
|  0|     7|               0|
|  1|     6|               1|
|  2|     9|               1|
|  3|     7|               0|
|  4|     3|               1|
|  5|     8|               0|
|  6|     9|               1|
|  7|     8|               0|
|  8|     3|               1|
|  9|     8|               0|
+---+------+----------------+

14. Как извлечь элементы по заданным позициям из столбца?

Уровень: L2

Входные данные:

from pyspark.sql.functions import rand

df = spark.range(10)
df = df.withColumn("random", ((rand(seed=42) * 10) + 1).cast("int"))

pos = [0, 4, 8, 5]
+---+------+
| id|random|
+---+------+
|  0|     7|
|  1|     6|
|  2|     9|
|  3|     7|
|  4|     3|
|  5|     8|
|  6|     9|
|  7|     8|
|  8|     3|
|  9|     8|
+---+------+
Решение
from pyspark.sql.window import Window
from pyspark.sql.functions import row_number, monotonically_increasing_id

pos = [0, 4, 8, 5]

w = Window.orderBy(monotonically_increasing_id())
df = df.withColumn("index", row_number().over(w) - 1)

# Фильтруем строки по указанным позициям
df_filtered = df.filter(df.index.isin(pos))
df_filtered.show()
+---+------+-----+
| id|random|index|
+---+------+-----+
|  0|     7|    0|
|  4|     3|    4|
|  5|     8|    5|
|  8|     3|    8|
+---+------+-----+

15. Как объединить два DataFrame вертикально?

Уровень: L1

Входные данные:

df_A = spark.createDataFrame(
    [("apple", 3, 5), ("banana", 1, 10), ("orange", 2, 8)],
    ["Name", "Col_1", "Col_2"])

df_B = spark.createDataFrame(
    [("apple", 3, 5), ("banana", 1, 15), ("grape", 4, 6)],
    ["Name", "Col_1", "Col_3"])

df_A.show()
df_B.show()
+------+-----+-----+
|  Name|Col_1|Col_2|
+------+-----+-----+
| apple|    3|    5|
|banana|    1|   10|
|orange|    2|    8|
+------+-----+-----+

+------+-----+-----+
|  Name|Col_1|Col_3|
+------+-----+-----+
| apple|    3|    5|
|banana|    1|   15|
| grape|    4|    6|
+------+-----+-----+
Решение
df_A.union(df_B).show()
+------+-----+-----+
|  Name|Col_1|Col_2|
+------+-----+-----+
| apple|    3|    5|
|banana|    1|   10|
|orange|    2|    8|
| apple|    3|    5|
|banana|    1|   15|
| grape|    4|    6|
+------+-----+-----+

16. Как вычислить среднеквадратичную ошибку (MSE)?

Уровень: L2

Входные данные:

data = [(1, 1), (2, 4), (3, 9), (4, 16), (5, 25)]
df = spark.createDataFrame(data, ["actual", "predicted"])

df.show()
+------+---------+
|actual|predicted|
+------+---------+
|     1|        1|
|     2|        4|
|     3|        9|
|     4|       16|
|     5|       25|
+------+---------+
Решение
from pyspark.sql.functions import col, pow

# Квадрат разности
df = df.withColumn("squared_error", pow((col("actual") - col("predicted")), 2))

# Среднее значение
mse = df.agg({"squared_error": "avg"}).collect()[0][0]

print(f"Mean Squared Error (MSE) = {mse}")
Mean Squared Error (MSE) = 116.8

17. Как привести первый символ каждого элемента к верхнему регистру?

Уровень: L1

Входные данные:

data = [("john",), ("alice",), ("bob",)]
df = spark.createDataFrame(data, ["name"])

df.show()
+-----+
| name|
+-----+
| john|
|alice|
|  bob|
+-----+
Решение
from pyspark.sql.functions import initcap

df = df.withColumn("name", initcap(df["name"]))
df.show()
+-----+
| name|
+-----+
| John|
|Alice|
|  Bob|
+-----+

18. Как вычислить сводную статистику по всем столбцам DataFrame?

Уровень: L1

Входные данные:

data = [('James',   34, 55000),
        ('Michael', 30, 70000),
        ('Robert',  37, 60000),
        ('Maria',   29, 80000),
        ('Jen',     32, 65000)]

df = spark.createDataFrame(data, ["name", "age", "salary"])
df.show()
+-------+---+------+
|   name|age|salary|
+-------+---+------+
|  James| 34| 55000|
|Michael| 30| 70000|
| Robert| 37| 60000|
|  Maria| 29| 80000|
|    Jen| 32| 65000|
+-------+---+------+
Решение
summary = df.summary()
summary.show()
+-------+-------+------------------+-----------------+
|summary|   name|               age|           salary|
+-------+-------+------------------+-----------------+
|  count|      5|                 5|                5|
|   mean|   null|              32.4|          66000.0|
| stddev|   null|3.2093613071762417|9617.692030835673|
|    min|  James|                29|            55000|
|    25%|   null|                30|            60000|
|    50%|   null|                32|            65000|
|    75%|   null|                34|            70000|
|    max| Robert|                37|            80000|
+-------+-------+------------------+-----------------+

19. Как посчитать количество символов в каждом слове столбца?

Уровень: L1

Входные данные:

data = [("john",), ("alice",), ("bob",)]
df = spark.createDataFrame(data, ["name"])
+-----+
| name|
+-----+
| john|
|alice|
|  bob|
+-----+
Решение
from pyspark.sql import functions as F

df = df.withColumn('word_length', F.length(df.name))
df.show()
+-----+-----------+
| name|word_length|
+-----+-----------+
| john|          4|
|alice|          5|
|  bob|          3|
+-----+-----------+

20. Как вычислить разницу разностей между последовательными числами столбца?

Уровень: L2

Входные данные:

data = [('James',   34, 55000),
        ('Michael', 30, 70000),
        ('Robert',  37, 60000),
        ('Maria',   29, 80000),
        ('Jen',     32, 65000)]

df = spark.createDataFrame(data, ["name", "age", "salary"])
+-------+---+------+
|   name|age|salary|
+-------+---+------+
|  James| 34| 55000|
|Michael| 30| 70000|
| Robert| 37| 60000|
|  Maria| 29| 80000|
|    Jen| 32| 65000|
+-------+---+------+
Решение
from pyspark.sql import functions as F
from pyspark.sql.window import Window

df = df.withColumn("id", F.monotonically_increasing_id())
window = Window.orderBy("id")

# Предыдущее значение (лаг)
df = df.withColumn("prev_value", F.lag(df.salary).over(window))

# Разность с лагом
df = df.withColumn("diff",
    F.when(F.isnull(df.salary - df.prev_value), 0)
     .otherwise(df.salary - df.prev_value)).drop("id")

df.show()
+-------+---+------+----------+------+
|   name|age|salary|prev_value|  diff|
+-------+---+------+----------+------+
|  James| 34| 55000|      null|     0|
|Michael| 30| 70000|     55000| 15000|
| Robert| 37| 60000|     70000|-10000|
|  Maria| 29| 80000|     60000| 20000|
|    Jen| 32| 65000|     80000|-15000|
+-------+---+------+----------+------+

21. Как получить день месяца, номер недели, день года и день недели из строки даты?

Уровень: L2

Входные данные:

data = [("2023-05-18", "01 Jan 2010"),
        ("2023-12-31", "01 Jan 2010")]
df = spark.createDataFrame(data, ["date_str_1", "date_str_2"])
+----------+-----------+
|date_str_1| date_str_2|
+----------+-----------+
|2023-05-18|01 Jan 2010|
|2023-12-31|01 Jan 2010|
+----------+-----------+
Решение
from pyspark.sql.functions import to_date, dayofmonth, weekofyear, dayofyear, dayofweek

df = df.withColumn("date_1", to_date(df.date_str_1, 'yyyy-MM-dd'))
df = df.withColumn("date_2", to_date(df.date_str_2, 'dd MMM yyyy'))

df = (df
    .withColumn("day_of_month", dayofmonth(df.date_1))
    .withColumn("week_number",  weekofyear(df.date_1))
    .withColumn("day_of_year",  dayofyear(df.date_1))
    .withColumn("day_of_week",  dayofweek(df.date_1)))

df.show()
+----------+-----------+----------+----------+------------+-----------+-----------+-----------+
|date_str_1| date_str_2|    date_1|    date_2|day_of_month|week_number|day_of_year|day_of_week|
+----------+-----------+----------+----------+------------+-----------+-----------+-----------+
|2023-05-18|01 Jan 2010|2023-05-18|2010-01-01|          18|         20|        138|          5|
|2023-12-31|01 Jan 2010|2023-12-31|2010-01-01|          31|         52|        365|          1|
+----------+-----------+----------+----------+------------+-----------+-----------+-----------+

22. Как преобразовать строку «год-месяц» в дату, соответствующую 4-му дню месяца?

Уровень: L2

Входные данные:

df = spark.createDataFrame(
    [('Jan 2010',), ('Feb 2011',), ('Mar 2012',)],
    ['MonthYear'])
+---------+
|MonthYear|
+---------+
| Jan 2010|
| Feb 2011|
| Mar 2012|
+---------+
Решение
from pyspark.sql.functions import expr

# Преобразуем строку в дату (по умолчанию 1-й день месяца)
df = df.withColumn('Date', expr("to_date(MonthYear, 'MMM yyyy')"))
df.show()

# Заменяем день на 4-й
df = df.withColumn('Date', expr("date_add(date_sub(Date, day(Date) - 1), 3)"))
df.show()
+---------+----------+
|MonthYear|      Date|
+---------+----------+
| Jan 2010|2010-01-01|
| Feb 2011|2011-02-01|
| Mar 2012|2012-03-01|
+---------+----------+

+---------+----------+
|MonthYear|      Date|
+---------+----------+
| Jan 2010|2010-01-04|
| Feb 2011|2011-02-04|
| Mar 2012|2012-03-04|
+---------+----------+

23. Как отфильтровать слова, содержащие не менее 2 гласных?

Уровень: L3

Входные данные:

df = spark.createDataFrame(
    [('Apple',), ('Orange',), ('Plan',), ('Python',), ('Money',)],
    ['Word'])

df.show()
+------+
|  Word|
+------+
| Apple|
|Orange|
|  Plan|
|Python|
| Money|
+------+
Решение
from pyspark.sql.functions import col, length, translate

# Фильтруем слова с не менее 2 гласными
df_filtered = df.where(
    (length(col('Word')) - length(translate(col('Word'), 'AEIOUaeiou', ''))) >= 2
)
df_filtered.show()
+------+
|  Word|
+------+
| Apple|
|Orange|
| Money|
+------+

24. Как отфильтровать корректные email-адреса из списка?

Уровень: L3

Входные данные:

from pyspark.sql import functions as F

data = ['buying books at amazom.com', 'rameses@egypt.com',
        'matt@t.co', 'narendra@modi.com']

df = spark.createDataFrame(data, "string")
df.show(truncate=False)
+--------------------------+
|value                     |
+--------------------------+
|buying books at amazom.com|
|rameses@egypt.com         |
|matt@t.co                 |
|narendra@modi.com         |
+--------------------------+
Решение
# Регулярное выражение для email
pattern = "^[a-zA-Z0-9_.+-]+@[a-zA-Z0-9-]+\\.[a-zA-Z0-9-.]+$"

df_filtered = df.filter(F.col("value").rlike(pattern))
df_filtered.show()
+-----------------+
|            value|
+-----------------+
|rameses@egypt.com|
|        matt@t.co|
|narendra@modi.com|
+-----------------+

25. Как выполнить Pivot для PySpark DataFrame?

Уровень: L3

Преобразуйте категории столбца region в отдельные столбцы и просуммируйте revenue.

Входные данные:

data = [
    (2021, 1, "US", 5000),
    (2021, 1, "EU", 4000),
    (2021, 2, "US", 5500),
    (2021, 2, "EU", 4500),
    (2021, 3, "US", 6000),
    (2021, 3, "EU", 5000),
    (2021, 4, "US", 7000),
    (2021, 4, "EU", 6000),
]

columns = ["year", "quarter", "region", "revenue"]
df = spark.createDataFrame(data, columns)
df.show()
+----+-------+------+-------+
|year|quarter|region|revenue|
+----+-------+------+-------+
|2021|      1|    US|   5000|
|2021|      1|    EU|   4000|
|2021|      2|    US|   5500|
|2021|      2|    EU|   4500|
|2021|      3|    US|   6000|
|2021|      3|    EU|   5000|
|2021|      4|    US|   7000|
|2021|      4|    EU|   6000|
+----+-------+------+-------+
Решение
pivot_df = df.groupBy("year", "quarter").pivot("region").sum("revenue")
pivot_df.show()
+----+-------+----+----+
|year|quarter|  EU|  US|
+----+-------+----+----+
|2021|      2|4500|5500|
|2021|      1|4000|5000|
|2021|      3|5000|6000|
|2021|      4|6000|7000|
+----+-------+----+----+

26. Как получить среднее значение переменной, сгруппированной по другой переменной?

Уровень: L3

Входные данные:

data = [("1001", "Laptop",     1000),
        ("1002", "Mouse",        50),
        ("1003", "Laptop",     1200),
        ("1004", "Mouse",        30),
        ("1005", "Smartphone", 700)]

columns = ["OrderID", "Product", "Price"]
df = spark.createDataFrame(data, columns)
df.show()
+-------+----------+-----+
|OrderID|   Product|Price|
+-------+----------+-----+
|   1001|    Laptop| 1000|
|   1002|     Mouse|   50|
|   1003|    Laptop| 1200|
|   1004|     Mouse|   30|
|   1005|Smartphone|  700|
+-------+----------+-----+
Решение
from pyspark.sql.functions import mean

result = df.groupBy("Product").agg(mean("Price").alias("Avg_Price"))
result.show()
+----------+---------+
|   Product|Avg_Price|
+----------+---------+
|    Laptop|   1100.0|
|     Mouse|     40.0|
|Smartphone|    700.0|
+----------+---------+

27. Как вычислить евклидово расстояние между двумя столбцами?

Уровень: L3

Вычислите евклидово расстояние между сериями series1 и series2 без использования готовой формулы.

Входные данные:

data = [(1, 10), (2, 9), (3, 8), (4, 7), (5, 6),
        (6, 5),  (7, 4), (8, 3), (9, 2), (10, 1)]

df = spark.createDataFrame(data, ["series1", "series2"])
df.show()
+-------+-------+
|series1|series2|
+-------+-------+
|      1|     10|
|      2|      9|
|      3|      8|
|      4|      7|
|      5|      6|
|      6|      5|
|      7|      4|
|      8|      3|
|      9|      2|
|     10|      1|
+-------+-------+
Решение
from pyspark.sql.functions import expr

# Квадрат разности
df = df.withColumn("squared_diff", expr("POW(series1 - series2, 2)"))

# Корень из суммы квадратов
df.agg(expr("SQRT(SUM(squared_diff))").alias("euclidean_distance")).show()
+------------------+
|euclidean_distance|
+------------------+
|18.16590212458495|
+------------------+

28. Как заменить пробелы в строке наименее частым символом?

Уровень: L3

Замените пробелы в строке символом с наименьшей частотой встречаемости.

Входные данные:

df = spark.createDataFrame([('dbc deb abed gade',)], ["string"])
df.show()
+-----------------+
|           string|
+-----------------+
|dbc deb abed gade|
+-----------------+

Ожидаемый результат:

+-----------------+-----------------+
|           string|  modified_string|
+-----------------+-----------------+
|dbc deb abed gade|dbccdebcabedcgade|
+-----------------+-----------------+
Решение
from pyspark.sql.functions import udf
from pyspark.sql.types import StringType
from collections import Counter

def least_freq_char_replace_spaces(s):
    counter = Counter(s.replace(" ", ""))
    least_freq_char = min(counter, key=counter.get)
    return s.replace(' ', least_freq_char)

udf_replace = udf(least_freq_char_replace_spaces, StringType())

df.withColumn('modified_string', udf_replace(df['string'])).show()
+-----------------+-----------------+
|           string|  modified_string|
+-----------------+-----------------+
|dbc deb abed gade|dbccdebcabedcgade|
+-----------------+-----------------+

29. Как создать временной ряд из 10 суббот начиная с '2000-01-01'?

Уровень: L3

Ожидаемый результат (значения случайные):

+----------+--------------+
|      date|random_numbers|
+----------+--------------+
|2000-01-01|             8|
|2000-01-08|             3|
|2000-01-15|             8|
|2000-01-22|             5|
|2000-01-29|             4|
|2000-02-05|             6|
|2000-02-12|             8|
|2000-02-19|             1|
|2000-02-26|             9|
|2000-03-04|             3|
+----------+--------------+
Решение
from pyspark.sql.functions import expr, explode, sequence, rand

start_date = '2000-01-01'
end_date   = '2000-03-04'

# Генерируем последовательность дат
df = spark.range(1).select(
    explode(
        sequence(
            expr(f"date '{start_date}'"),
            expr(f"date '{end_date}'"),
            expr("interval 1 day")
        )
    ).alias("date")
)

# Оставляем только субботы (7 в Spark)
df = df.filter(expr("dayofweek(date) = 7"))

# Добавляем случайные числа
df = df.withColumn("random_numbers", ((rand(seed=42) * 10) + 1).cast("int"))

df.show()
+----------+--------------+
|      date|random_numbers|
+----------+--------------+
|2000-01-01|             8|
|2000-01-08|             3|
|2000-01-15|             8|
|2000-01-22|             5|
|2000-01-29|             4|
|2000-02-05|             6|
|2000-02-12|             8|
|2000-02-19|             1|
|2000-02-26|             9|
|2000-03-04|             3|
+----------+--------------+

30. Как загрузить CSV-файл и получить форму DataFrame?

Уровень: L1

Загрузите датасет Churn Modelling и выведите количество строк, столбцов и типы данных.

Входные данные:

from pyspark import SparkFiles

url = "https://raw.githubusercontent.com/selva86/datasets/master/Churn_Modelling.csv"
spark.sparkContext.addFile(url)

df = spark.read.csv(SparkFiles.get("Churn_Modelling.csv"),
                    header=True, inferSchema=True)

df.show(5, truncate=False)
+---------+----------+--------+-----------+---------+------+---+------+---------+-------------+---------+--------------+---------------+------+
|RowNumber|CustomerId|Surname |CreditScore|Geography|Gender|Age|Tenure|Balance  |NumOfProducts|HasCrCard|IsActiveMember|EstimatedSalary|Exited|
+---------+----------+--------+-----------+---------+------+---+------+---------+-------------+---------+--------------+---------------+------+
|1        |15634602  |Hargrave|619        |France   |Female|42 |2     |0.0      |1            |1        |1             |101348.88      |1     |
|2        |15647311  |Hill    |608        |Spain    |Female|41 |1     |83807.86 |1            |0        |1             |112542.58      |0     |
+---------+----------+--------+-----------+---------+------+---+------+---------+-------------+---------+--------------+---------------+------+
only showing top 5 rows
Решение
# Количество строк
nrows = df.count()
print("Количество строк:", nrows)

# Количество столбцов
ncols = len(df.columns)
print("Количество столбцов:", ncols)

# Типы данных каждого столбца
datatypes = df.dtypes
print("Типы данных:", datatypes)
Количество строк:  10000
Количество столбцов:  14
Типы данных:  [('RowNumber', 'int'), ('CustomerId', 'int'), ('Surname', 'string'), ('CreditScore', 'int'), ('Geography', 'string'), ('Gender', 'string'), ('Age', 'int'), ('Tenure', 'int'), ('Balance', 'double'), ('NumOfProducts', 'int'), ('HasCrCard', 'int'), ('IsActiveMember', 'int'), ('EstimatedSalary', 'double'), ('Exited', 'int')]

31. Как переименовать конкретные столбцы DataFrame?

Уровень: L2

Входные данные:

df = spark.createDataFrame(
    [('Alice', 1, 30), ('Bob', 2, 35)],
    ["name", "age", "qty"])

df.show()

old_names = ["qty", "age"]
new_names = ["user_qty", "user_age"]
+-----+---+---+
| name|age|qty|
+-----+---+---+
|Alice|  1| 30|
|  Bob|  2| 35|
+-----+---+---+
Решение
for old_name, new_name in zip(old_names, new_names):
    df = df.withColumnRenamed(old_name, new_name)

df.show()
+-----+--------+--------+
| name|user_age|user_qty|
+-----+--------+--------+
|Alice|       1|      30|
|  Bob|       2|      35|
+-----+--------+--------+

32. Как проверить наличие пропущенных значений и подсчитать их количество?

Уровень: L2

Входные данные:

df = spark.createDataFrame([
    ("A",  1,    None),
    ("B",  None, "123"),
    ("B",  3,    "456"),
    ("D",  None, None),
], ["Name", "Value", "id"])

df.show()
+----+-----+----+
|Name|Value|  id|
+----+-----+----+
|   A|    1|null|
|   B| null| 123|
|   B|    3| 456|
|   D| null|null|
+----+-----+----+
Решение
from pyspark.sql.functions import col, sum

missing = df.select(*(sum(col(c).isNull().cast("int")).alias(c) for c in df.columns))
has_missing = any(missing.collect()[0].asDict().values())
print("Есть пропуски:", has_missing)

missing_count = missing.collect()[0].asDict()
print("Количество пропусков:", missing_count)
Есть пропуски: True
Количество пропусков: {'Name': 0, 'Value': 2, 'id': 2}

33. Как заменить пропущенные значения в числовых столбцах средним?

Уровень: L2

Входные данные:

df = spark.createDataFrame([
    ("A",  1,    None),
    ("B",  None, 123),
    ("B",  3,    456),
    ("D",  6,    None),
], ["Name", "var1", "var2"])

df.show()
+----+----+----+
|Name|var1|var2|
+----+----+----+
|   A|   1|null|
|   B|null| 123|
|   B|   3| 456|
|   D|   6|null|
+----+----+----+
Решение
from pyspark.ml.feature import Imputer

column_names = ["var1", "var2"]

imputer = Imputer(inputCols=column_names, outputCols=column_names, strategy="mean")
model = imputer.fit(df)
imputed_df = model.transform(df)

imputed_df.show()
+----+----+----+
|Name|var1|var2|
+----+----+----+
|   A|   1| 289|
|   B|   3| 123|
|   B|   3| 456|
|   D|   6| 289|
+----+----+----+

34. Как изменить порядок столбцов DataFrame?

Уровень: L1

Входные данные:

data = [("John", "Doe", 30), ("Jane", "Doe", 25), ("Alice", "Smith", 22)]
df = spark.createDataFrame(data, ["First_Name", "Last_Name", "Age"])

df.show()
+----------+---------+---+
|First_Name|Last_Name|Age|
+----------+---------+---+
|      John|      Doe| 30|
|      Jane|      Doe| 25|
|     Alice|    Smith| 22|
+----------+---------+---+
Решение
new_order = ["Age", "First_Name", "Last_Name"]
df = df.select(*new_order)
df.show()
+---+----------+---------+
|Age|First_Name|Last_Name|
+---+----------+---------+
| 30|      John|      Doe|
| 25|      Jane|      Doe|
| 22|     Alice|    Smith|
+---+----------+---------+

35. Как убрать научную нотацию в PySpark DataFrame?

Входные данные:

df = spark.createDataFrame(
    [(1, 0.000000123), (2, 0.000023456), (3, 0.000345678)],
    ["id", "your_column"])

df.show()
+---+-----------+
| id|your_column|
+---+-----------+
|  1|   1.23E-7|
|  2| 2.3456E-5|
|  3|3.45678E-4|
+---+-----------+
Решение
from pyspark.sql.functions import format_number

decimal_places = 10
df = df.withColumn("your_column", format_number("your_column", decimal_places))
df.show()
+---+------------+
| id| your_column|
+---+------------+
|  1|0.0000001230|
|  2|0.0000234560|
|  3|0.0003456780|
+---+------------+

36. Как представить все значения DataFrame в виде процентов?

Уровень: L2

Входные данные:

data = [(0.1, .08), (0.2, .06), (0.33, .02)]
df = spark.createDataFrame(data, ["numbers_1", "numbers_2"])
+---------+---------+
|numbers_1|numbers_2|
+---------+---------+
|      0.1|     0.08|
|      0.2|     0.06|
|     0.33|     0.02|
+---------+---------+
Решение
from pyspark.sql.functions import concat, col, lit

columns = ["numbers_1", "numbers_2"]

for col_name in columns:
    df = df.withColumn(col_name,
        concat((col(col_name) * 100).cast('decimal(10, 2)'), lit("%")))

df.show()
+---------+---------+
|numbers_1|numbers_2|
+---------+---------+
|  10.00% |   8.00%|
|  20.00% |   6.00%|
|  33.00% |   2.00%|
+---------+---------+

37. Как отфильтровать каждую n-ю строку DataFrame?

Уровень: L2

Входные данные:

from pyspark.sql.window import Window
from pyspark.sql.functions import row_number, monotonically_increasing_id

data = [("Alice", 1), ("Bob", 2), ("Charlie", 3), ("Dave", 4), ("Eve", 5),
        ("Frank", 6), ("Grace", 7), ("Hannah", 8), ("Igor", 9), ("Jack", 10)]

df = spark.createDataFrame(data, ["Name", "Number"])
+-------+------+
|   Name|Number|
+-------+------+
|  Alice|     1|
|    Bob|     2|
|Charlie|     3|
|   Dave|     4|
|    Eve|     5|
|  Frank|     6|
|  Grace|     7|
| Hannah|     8|
|   Igor|     9|
|   Jack|    10|
+-------+------+
Решение
window = Window.orderBy(monotonically_increasing_id())
df = df.withColumn("rn", row_number().over(window))

n = 5  # каждая 5-я строка
df = df.filter((df.rn % n) == 0)
df.show()
+----+------+---+
|Name|Number| rn|
+----+------+---+
| Eve|     5|  5|
|Jack|    10| 10|
+----+------+---+

38. Как получить номер строки с n-м наибольшим значением в столбце?

Уровень: L2

Входные данные:

from pyspark.sql import Row

data = [
    Row(id=1, column1=5),
    Row(id=2, column1=8),
    Row(id=3, column1=12),
    Row(id=4, column1=1),
    Row(id=5, column1=15),
    Row(id=6, column1=7),
]

df = spark.createDataFrame(data)
+---+-------+
| id|column1|
+---+-------+
|  1|      5|
|  2|      8|
|  3|     12|
|  4|      1|
|  5|     15|
|  6|      7|
+---+-------+
Решение
from pyspark.sql.window import Window
from pyspark.sql.functions import desc, row_number

window = Window.orderBy(desc("column1"))
df = df.withColumn("row_number", row_number().over(window))

n = 3  # 3-е наибольшее значение
row = df.filter(df.row_number == n).first()

if row:
    print("Номер строки:", row.row_number)
    print("Значение в столбце:", row.column1)
Номер строки: 3
Значение в столбце: 8

39. Как получить последние n строк, где сумма строки > 100?

Уровень: L2

Входные данные:

data = [(10, 25, 70),
        (40,  5, 20),
        (70, 80, 100),
        (10,  2, 60),
        (40, 50, 20)]

df = spark.createDataFrame(data, ["col1", "col2", "col3"])
+----+----+----+
|col1|col2|col3|
+----+----+----+
|  10|  25|  70|
|  40|   5|  20|
|  70|  80| 100|
|  10|   2|  60|
|  40|  50|  20|
+----+----+----+
Решение
from pyspark.sql import functions as F
from functools import reduce

# Добавляем столбец суммы строки
df = df.withColumn('row_sum', reduce(lambda a, b: a + b, [F.col(x) for x in df.columns]))

# Фильтруем строки, где сумма > 100
df = df.filter(F.col('row_sum') > 100)

# Берём последние 2 строки
df = df.withColumn('id', F.monotonically_increasing_id())
df_last_2 = df.sort(F.desc('id')).limit(2)

df_last_2.show()
+----+----+----+-------+-----------+
|col1|col2|col3|row_sum|         id|
+----+----+----+-------+-----------+
|  40|  50|  20|    110|25769803776|
|  70|  80| 100|    250|17179869184|
+----+----+----+-------+-----------+

40. Как создать столбец с отношением min/max в каждой строке?

Уровень: L2

Входные данные:

data = [(1, 2, 3), (4, 5, 6), (7, 8, 9), (10, 11, 12)]
df = spark.createDataFrame(data, ["col1", "col2", "col3"])
+----+----+----+
|col1|col2|col3|
+----+----+----+
|   1|   2|   3|
|   4|   5|   6|
|   7|   8|   9|
|  10|  11|  12|
+----+----+----+
Решение
from pyspark.sql.functions import udf, array
from pyspark.sql.types import FloatType

def min_max_ratio(row):
    return float(min(row)) / max(row)

min_max_ratio_udf = udf(min_max_ratio, FloatType())

df = df.withColumn('min_by_max', min_max_ratio_udf(array(df.columns)))
df.show()
+----+----+----+----------+
|col1|col2|col3|min_by_max|
+----+----+----+----------+
|   1|   2|   3|0.33333334|
|   4|   5|   6| 0.6666667|
|   7|   8|   9| 0.7777778|
|  10|  11|  12| 0.8333333|
+----+----+----+----------+

41. Как создать столбец с предпоследним значением каждой строки?

Уровень: L2

Создайте новый столбец Penultimate, содержащий второе по величине значение в каждой строке.

Входные данные:

data = [(10, 20, 30),
        (40, 60, 50),
        (80, 70, 90)]

df = spark.createDataFrame(data, ["Column1", "Column2", "Column3"])
+-------+-------+-------+
|Column1|Column2|Column3|
+-------+-------+-------+
|     10|     20|     30|
|     40|     60|     50|
|     80|     70|     90|
+-------+-------+-------+
Решение
from pyspark.sql import functions as F
from pyspark.sql.types import ArrayType, IntegerType

sort_array_asc = F.udf(lambda arr: sorted(arr), ArrayType(IntegerType()))

df = df.withColumn("row_as_array", sort_array_asc(F.array(df.columns)))
df = df.withColumn("Penultimate", df['row_as_array'].getItem(1))
df = df.drop('row_as_array')

df.show()
+-------+-------+-------+-----------+
|Column1|Column2|Column3|Penultimate|
+-------+-------+-------+-----------+
|     10|     20|     30|         20|
|     40|     60|     50|         50|
|     80|     70|     90|         80|
+-------+-------+-------+-----------+

42. Как нормализовать все столбцы DataFrame?

Уровень: L2

Нормализуйте все столбцы: вычтите среднее и разделите на стандартное отклонение.

Входные данные:

data = [(1, 2, 3),
        (2, 3, 4),
        (3, 4, 5),
        (4, 5, 6)]

df = spark.createDataFrame(data, ["Col1", "Col2", "Col3"])
+----+----+----+
|Col1|Col2|Col3|
+----+----+----+
|   1|   2|   3|
|   2|   3|   4|
|   3|   4|   5|
|   4|   5|   6|
+----+----+----+
Решение
from pyspark.ml.feature import VectorAssembler, StandardScaler

input_cols = ["Col1", "Col2", "Col3"]

assembler = VectorAssembler(inputCols=input_cols, outputCol="features")
df_assembled = assembler.transform(df)

scaler = StandardScaler(inputCol="features", outputCol="scaled_features",
                        withStd=True, withMean=True)
scalerModel = scaler.fit(df_assembled)
df_normalized = scalerModel.transform(df_assembled).drop('features')

df_normalized.show(truncate=False)
+----+----+----+-------------------------------------------------------------+
|Col1|Col2|Col3|scaled_features                                              |
+----+----+----+-------------------------------------------------------------+
|1   |2   |3   |[-1.161895003862225,-1.161895003862225,-1.161895003862225]   |
|2   |3   |4   |[-0.387298334620..., -0.387298334620..., -0.387298334620...] |
|3   |4   |5   |[0.387298334620..., 0.387298334620..., 0.387298334620...]    |
|4   |5   |6   |[1.161895003862225,1.161895003862225,1.161895003862225]      |
+----+----+----+-------------------------------------------------------------+

43. Как найти позиции, где значения двух столбцов совпадают?

Уровень: L1

Входные данные:

data = [("John", "John"), ("Lily", "Lucy"), ("Sam", "Sam"), ("Lucy", "Lily")]
df = spark.createDataFrame(data, ["Name1", "Name2"])
+-----+-----+
|Name1|Name2|
+-----+-----+
| John| John|
| Lily| Lucy|
|  Sam|  Sam|
| Lucy| Lily|
+-----+-----+
Решение
from pyspark.sql.functions import when, col

df = df.withColumn("Match",
    when(col("Name1") == col("Name2"), True).otherwise(False))

df.show()
+-----+-----+-----+
|Name1|Name2|Match|
+-----+-----+-----+
| John| John| true|
| Lily| Lucy|false|
|  Sam|  Sam| true|
| Lucy| Lily|false|
+-----+-----+-----+

44. Как создать лаги и лиды столбца по группе?

Уровень: L2

Входные данные:

data = [
    ("2023-01-01", "Store1", 100), ("2023-01-02", "Store1", 150),
    ("2023-01-03", "Store1", 200), ("2023-01-04", "Store1", 250),
    ("2023-01-05", "Store1", 300), ("2023-01-01", "Store2", 50),
    ("2023-01-02", "Store2", 60),  ("2023-01-03", "Store2", 80),
    ("2023-01-04", "Store2", 90),  ("2023-01-05", "Store2", 120),
]

df = spark.createDataFrame(data, ["Date", "Store", "Sales"])
+----------+------+-----+
|      Date| Store|Sales|
+----------+------+-----+
|2023-01-01|Store1|  100|
|2023-01-02|Store1|  150|
|2023-01-03|Store1|  200|
|2023-01-04|Store1|  250|
|2023-01-05|Store1|  300|
|2023-01-01|Store2|   50|
|2023-01-02|Store2|   60|
|2023-01-03|Store2|   80|
|2023-01-04|Store2|   90|
|2023-01-05|Store2|  120|
+----------+------+-----+
Решение
from pyspark.sql.functions import lag, lead, to_date
from pyspark.sql.window import Window

df = df.withColumn("Date", to_date(df.Date, 'yyyy-MM-dd'))

windowSpec = Window.partitionBy("Store").orderBy("Date")

df = df.withColumn("Lag_Sales",  lag(df["Sales"]).over(windowSpec))
df = df.withColumn("Lead_Sales", lead(df["Sales"]).over(windowSpec))

df.show()
+----------+------+-----+---------+----------+
|      Date| Store|Sales|Lag_Sales|Lead_Sales|
+----------+------+-----+---------+----------+
|2023-01-01|Store1|  100|     null|       150|
|2023-01-02|Store1|  150|      100|       200|
|2023-01-03|Store1|  200|      150|       250|
|2023-01-04|Store1|  250|      200|       300|
|2023-01-05|Store1|  300|      250|      null|
|2023-01-01|Store2|   50|     null|        60|
|2023-01-02|Store2|   60|       50|        80|
|2023-01-03|Store2|   80|       60|        90|
|2023-01-04|Store2|   90|       80|       120|
|2023-01-05|Store2|  120|       90|      null|
+----------+------+-----+---------+----------+

45. Как получить частоту уникальных значений во всём DataFrame?

Уровень: L3

Входные данные:

data = [(1, 2, 3), (2, 3, 4), (1, 2, 3), (4, 5, 6), (2, 3, 4)]
df = spark.createDataFrame(data, ["Column1", "Column2", "Column3"])
+-------+-------+-------+
|Column1|Column2|Column3|
+-------+-------+-------+
|      1|      2|      3|
|      2|      3|      4|
|      1|      2|      3|
|      4|      5|      6|
|      2|      3|      4|
+-------+-------+-------+
Решение
from pyspark.sql.functions import col

columns = df.columns

# Складываем все столбцы в один
df_single = None
for c in columns:
    tmp = df.select(col(c).alias("single_column"))
    df_single = tmp if df_single is None else df_single.union(tmp)

frequency_table = df_single.groupBy("single_column").count().orderBy('count', ascending=False)
frequency_table.show()
+-------------+-----+
|single_column|count|
+-------------+-----+
|            3|    4|
|            2|    4|
|            4|    3|
|            1|    2|
|            5|    1|
|            6|    1|
+-------------+-----+

46. Как заменить оба диагональных значения DataFrame нулями?

Уровень: L3

Замените значения на обеих диагоналях DataFrame нулями.

Входные данные:

data = [(1, 2, 3, 4), (2, 3, 4, 5), (1, 2, 3, 4), (4, 5, 6, 7)]
df = spark.createDataFrame(data, ["col_1", "col_2", "col_3", "col_4"])

df.show()
+-----+-----+-----+-----+
|col_1|col_2|col_3|col_4|
+-----+-----+-----+-----+
|    1|    2|    3|    4|
|    2|    3|    4|    5|
|    1|    2|    3|    4|
|    4|    5|    6|    7|
+-----+-----+-----+-----+
Решение
from pyspark.sql import functions as F
from pyspark.sql.window import Window
from pyspark.sql.functions import row_number, monotonically_increasing_id, when, col

w = Window.orderBy(monotonically_increasing_id())
df = df.withColumn("id", row_number().over(w) - 1)

# Основная диагональ
df = df.select([when(col("id") == i, 0).otherwise(col("col_" + str(i+1))).alias("col_" + str(i+1)) for i in range(4)])

# Обратная диагональ
df = df.withColumn("id", row_number().over(w) - 1)
df = df.withColumn("id_2", df.count() - 1 - df["id"])

df_with_diag_zero = df.select([
    when(col("id_2") == i, 0).otherwise(col("col_" + str(i+1))).alias("col_" + str(i+1))
    for i in range(4)])

df_with_diag_zero.show()
+-----+-----+-----+-----+
|col_1|col_2|col_3|col_4|
+-----+-----+-----+-----+
|    0|    2|    3|    0|
|    2|    0|    0|    5|
|    1|    0|    0|    4|
|    0|    5|    6|    0|
+-----+-----+-----+-----+

47. Как перевернуть строки DataFrame?

Уровень: L2

Входные данные:

data = [(1, 2, 3, 4), (2, 3, 4, 5), (3, 4, 5, 6), (4, 5, 6, 7)]
df = spark.createDataFrame(data, ["col_1", "col_2", "col_3", "col_4"])
+-----+-----+-----+-----+
|col_1|col_2|col_3|col_4|
+-----+-----+-----+-----+
|    1|    2|    3|    4|
|    2|    3|    4|    5|
|    3|    4|    5|    6|
|    4|    5|    6|    7|
+-----+-----+-----+-----+
Решение
from pyspark.sql.window import Window
from pyspark.sql.functions import row_number, monotonically_increasing_id

w = Window.orderBy(monotonically_increasing_id())
df = df.withColumn("id", row_number().over(w) - 1)

df_reversed = df.orderBy("id", ascending=False).drop("id")
df_reversed.show()
+-----+-----+-----+-----+
|col_1|col_2|col_3|col_4|
+-----+-----+-----+-----+
|    4|    5|    6|    7|
|    3|    4|    5|    6|
|    2|    3|    4|    5|
|    1|    2|    3|    4|
+-----+-----+-----+-----+

48. Как создать one-hot кодирование категориальной переменной?

Уровень: L2

Входные данные:

data = [("A", 10), ("A", 20), ("B", 30), ("B", 20),
        ("B", 30), ("C", 40), ("C", 10), ("D", 10)]
columns = ["Categories", "Value"]

df = spark.createDataFrame(data, columns)
df.show()
+----------+-----+
|Categories|Value|
+----------+-----+
|         A|   10|
|         A|   20|
|         B|   30|
|         B|   20|
|         B|   30|
|         C|   40|
|         C|   10|
|         D|   10|
+----------+-----+
Решение
from pyspark.ml.feature import StringIndexer, OneHotEncoder

# Шаг 1: StringIndexer — преобразует строки в числовые индексы
indexer = StringIndexer(inputCol="Categories", outputCol="Categories_Indexed")
indexerModel = indexer.fit(df)
indexed_df = indexerModel.transform(df)

# Шаг 2: OneHotEncoder
encoder = OneHotEncoder(inputCol="Categories_Indexed", outputCol="Categories_onehot")
encoded_df = encoder.fit(indexed_df).transform(indexed_df)
encoded_df = encoded_df.drop("Categories_Indexed")
encoded_df.show(truncate=False)
+----------+-----+-----------------+
|Categories|Value|Categories_onehot|
+----------+-----+-----------------+
|A         |10   |(3,[1],[1.0])    |
|A         |20   |(3,[1],[1.0])    |
|B         |30   |(3,[0],[1.0])    |
|B         |20   |(3,[0],[1.0])    |
|B         |30   |(3,[0],[1.0])    |
|C         |40   |(3,[2],[1.0])    |
|C         |10   |(3,[2],[1.0])    |
|D         |10   |(3,[],[])        |
+----------+-----+-----------------+

49. Как выполнить Pivot DataFrame (преобразование строк в столбцы)?

Уровень: L2

Преобразуйте категории столбца region в отдельные столбцы.

Входные данные:

data = [
    (2021, 1, "US", 5000), (2021, 1, "EU", 4000),
    (2021, 2, "US", 5500), (2021, 2, "EU", 4500),
    (2021, 3, "US", 6000), (2021, 3, "EU", 5000),
    (2021, 4, "US", 7000), (2021, 4, "EU", 6000),
]

columns = ["year", "quarter", "region", "revenue"]
df = spark.createDataFrame(data, columns)
Решение
pivot_df = df.groupBy("year", "quarter").pivot("region").sum("revenue")
pivot_df.show()
+----+-------+----+----+
|year|quarter|  EU|  US|
+----+-------+----+----+
|2021|      2|4500|5500|
|2021|      1|4000|5000|
|2021|      3|5000|6000|
|2021|      4|6000|7000|
+----+-------+----+----+

50. Как выполнить UnPivot DataFrame (преобразование столбцов в строки)?

Уровень: L2

Преобразуйте столбцы EU и US обратно в строки с колонками region и revenue.

Входные данные:

data = [(2021, 2, 4500, 5500),
        (2021, 1, 4000, 5000),
        (2021, 3, 5000, 6000),
        (2021, 4, 6000, 7000)]

columns = ["year", "quarter", "EU", "US"]
pivot_df = spark.createDataFrame(data, columns)
pivot_df.show()
+----+-------+----+----+
|year|quarter|  EU|  US|
+----+-------+----+----+
|2021|      2|4500|5500|
|2021|      1|4000|5000|
|2021|      3|5000|6000|
|2021|      4|6000|7000|
+----+-------+----+----+

Ожидаемый результат:

+----+-------+------+-------+
|year|quarter|region|revenue|
+----+-------+------+-------+
|2021|      2|    EU|   4500|
|2021|      2|    US|   5500|
|2021|      1|    EU|   4000|
|2021|      1|    US|   5000|
|2021|      3|    EU|   5000|
|2021|      3|    US|   6000|
|2021|      4|    EU|   6000|
|2021|      4|    US|   7000|
+----+-------+------+-------+
Решение
from pyspark.sql.functions import expr

unpivotExpr = "stack(2, 'EU', EU, 'US', US) as (region, revenue)"

unPivotDF = pivot_df.select("year", "quarter", expr(unpivotExpr)).where("revenue is not null")
unPivotDF.show()
+----+-------+------+-------+
|year|quarter|region|revenue|
+----+-------+------+-------+
|2021|      2|    EU|   4500|
|2021|      2|    US|   5500|
|2021|      1|    EU|   4000|
|2021|      1|    US|   5000|
|2021|      3|    EU|   5000|
|2021|      3|    US|   6000|
|2021|      4|    EU|   6000|
|2021|      4|    US|   7000|
+----+-------+------+-------+

51. Как заменить значения NA в DataFrame нулём?

Уровень: L1

Входные данные:

df = spark.createDataFrame(
    [(1, None), (None, 2), (3, 4), (5, None)],
    ["a", "b"])

df.show()
+----+----+
|   a|   b|
+----+----+
|   1|null|
|null|   2|
|   3|   4|
|   5|null|
+----+----+
Решение
df_imputed = df.fillna(0)
df_imputed.show()
+---+---+
|  a|  b|
+---+---+
|  1|  0|
|  0|  2|
|  3|  4|
|  5|  0|
+---+---+

52. Как определить непрерывные переменные в DataFrame?

Уровень: L3

Определите непрерывные переменные и создайте список их имён.

Входные данные (Churn Modelling датасет):

from pyspark import SparkFiles

url = "https://raw.githubusercontent.com/selva86/datasets/master/Churn_Modelling_m.csv"
spark.sparkContext.addFile(url)

df = spark.read.csv(SparkFiles.get("Churn_Modelling_m.csv"),
                    header=True, inferSchema=True)
df.show(2, truncate=False)
+---------+----------+--------+-----------+---------+------+---+------+--------+-------------+---------+--------------+---------------+------+
|RowNumber|CustomerId|Surname |CreditScore|Geography|Gender|Age|Tenure|Balance |NumOfProducts|HasCrCard|IsActiveMember|EstimatedSalary|Exited|
+---------+----------+--------+-----------+---------+------+---+------+--------+-------------+---------+--------------+---------------+------+
|1        |15634602  |Hargrave|619        |France   |Female|42 |2     |0.0     |1            |1        |1             |101348.88      |1     |
|2        |15647311  |Hill    |608        |Spain    |Female|41 |1     |83807.86|1            |0        |1             |112542.58      |0     |
+---------+----------+--------+-----------+---------+------+---+------+--------+-------------+---------+--------------+---------------+------+
only showing top 2 rows
Решение
from pyspark.sql.types import IntegerType, NumericType
from pyspark.sql.functions import approxCountDistinct

def detect_continuous_variables(df, distinct_threshold):
    """
    Определяет непрерывные переменные в PySpark DataFrame.
    :param distinct_threshold: порог — число уникальных значений > порога
    """
    continuous_columns = []
    for column in df.columns:
        dtype = df.schema[column].dataType
        if isinstance(dtype, (IntegerType, NumericType)):
            distinct_count = df.select(approxCountDistinct(column)).collect()[0][0]
            if distinct_count > distinct_threshold:
                continuous_columns.append(column)
    return continuous_columns

continuous_variables = detect_continuous_variables(df, 10)
print(continuous_variables)
['RowNumber', 'CustomerId', 'CreditScore', 'Age', 'Tenure', 'Balance', 'EstimatedSalary']

53. Как вычислить моду столбца PySpark DataFrame?

Уровень: L1

Входные данные:

data = [(1, 2, 3), (2, 2, 3), (2, 2, 4), (1, 2, 3), (1, 1, 3)]
columns = ["col1", "col2", "col3"]

df = spark.createDataFrame(data, columns)
df.show()
+----+----+----+
|col1|col2|col3|
+----+----+----+
|   1|   2|   3|
|   2|   2|   3|
|   2|   2|   4|
|   1|   2|   3|
|   1|   1|   3|
+----+----+----+
Решение
from pyspark.sql.functions import col

df_grouped = df.groupBy('col2').count()
mode_df = df_grouped.orderBy(col('count').desc()).limit(1)
mode_df.show()
+----+-----+
|col2|count|
+----+-----+
|   2|    4|
+----+-----+

54. Как найти путь установки Apache Spark и PySpark?

Уровень: L1

import findspark
findspark.init()

print(findspark.find())

import os
import pyspark

print(os.path.dirname(pyspark.__file__))
C:\spark\spark-3.3.2-bin-hadoop2
C:\spark\spark-3.3.2-bin-hadoop2\python\pyspark

55. Как привести значения столбца к нижнему регистру с помощью UDF?

Уровень: L2

Входные данные:

data = [('John Doe',    'NEW YORK'),
        ('Jane Doe',    'LOS ANGELES'),
        ('Mike Johnson', 'CHICAGO'),
        ('Sara Smith',  'SAN FRANCISCO')]

df = spark.createDataFrame(data, ['Name', 'City'])
df.show()
+------------+-------------+
|        Name|         City|
+------------+-------------+
|    John Doe|     NEW YORK|
|    Jane Doe|  LOS ANGELES|
|Mike Johnson|      CHICAGO|
|  Sara Smith|SAN FRANCISCO|
+------------+-------------+
Решение
from pyspark.sql.functions import udf
from pyspark.sql.types import StringType

def to_lower(s):
    if s is not None:
        return s.lower()

udf_to_lower = udf(to_lower, StringType())

df = df.withColumn('City_lower', udf_to_lower(df['City']))
df.show()
+------------+-------------+-------------+
|        Name|         City|   City_lower|
+------------+-------------+-------------+
|    John Doe|     NEW YORK|     new york|
|    Jane Doe|  LOS ANGELES|  los angeles|
|Mike Johnson|      CHICAGO|      chicago|
|  Sara Smith|SAN FRANCISCO|san francisco|
+------------+-------------+-------------+

56. Как преобразовать PySpark DataFrame в pandas DataFrame?

Уровень: L1

Входные данные:

data = [('John Doe',    'NEW YORK'),
        ('Jane Doe',    'LOS ANGELES'),
        ('Mike Johnson', 'CHICAGO'),
        ('Sara Smith',  'SAN FRANCISCO')]

pysparkDF = spark.createDataFrame(data, ['Name', 'City'])
pysparkDF.show()
+------------+-------------+
|        Name|         City|
+------------+-------------+
|    John Doe|     NEW YORK|
|    Jane Doe|  LOS ANGELES|
|Mike Johnson|      CHICAGO|
|  Sara Smith|SAN FRANCISCO|
+------------+-------------+
Решение
pandasDF = pysparkDF.toPandas()
print(pandasDF)
           Name           City
0      John Doe       NEW YORK
1      Jane Doe    LOS ANGELES
2  Mike Johnson        CHICAGO
3    Sara Smith  SAN FRANCISCO

57. Как просмотреть URL веб-интерфейса кластера PySpark?

Уровень: L1

print(spark.sparkContext.uiWebUrl)
http://DESKTOP-UL3QT3E.mshome.net:4040

58. Как просмотреть параметры конфигурации кластера PySpark?

Уровень: L1

for k, v in spark.sparkContext.getConf().getAll():
    print(f"{k} : {v}")
spark.app.name : PySpark 101 Exercises
spark.master : local[*]
spark.rdd.compress : True
spark.executor.id : driver
...

59. Как ограничить количество ядер, используемых PySpark?

Уровень: L1

from pyspark import SparkConf, SparkContext

conf = SparkConf()
conf.set("spark.executor.cores", "2")  # задаём нужное количество ядер
sc = SparkContext(conf=conf)

60. Как кэшировать PySpark DataFrame и очистить кэш?

Уровень: L2

В PySpark кэширование ускоряет повторный доступ к данным при итеративных вычислениях.

# Кэшируем DataFrame
df.cache()

# Очищаем кэш
df.unpersist()
DataFrame[Name: string, City: string, City_lower: string]

61. Как разделить PySpark DataFrame случайным образом в заданном соотношении?

Уровень: L1

# Разделяем данные в соотношении 0.8 / 0.2
train_data, test_data = df.randomSplit([0.8, 0.2], seed=42)

62. Как построить логистическую регрессию в PySpark?

Уровень: L2

Входные данные:

data = spark.createDataFrame([
    (0,  1.0, -1.0), (1,  2.0,  1.0), (1,  3.0, -2.0), (0,  4.0,  1.0),
    (1,  5.0, -3.0), (0,  6.0,  2.0), (1,  7.0, -1.0), (0,  8.0,  3.0),
    (1,  9.0, -2.0), (0, 10.0,  2.0), (1, 11.0, -3.0), (0, 12.0,  1.0),
    (1, 13.0, -1.0), (0, 14.0,  2.0), (1, 15.0, -2.0), (0, 16.0,  3.0),
    (1, 17.0, -3.0), (0, 18.0,  1.0), (1, 19.0, -1.0), (0, 20.0,  2.0)
], ["label", "feat1", "feat2"])
Решение
from pyspark.ml.feature import VectorAssembler
from pyspark.ml.classification import LogisticRegression

vecAssembler = VectorAssembler(inputCols=['feat1', 'feat2'], outputCol="features")
data = vecAssembler.transform(data)

lr = LogisticRegression(featuresCol='features', labelCol='label')
lr_model = lr.fit(data)

print(f"Коэффициенты: {lr_model.coefficients}")
print(f"Свободный член: {lr_model.intercept}")
Коэффициенты: [0.020277740475786673,-1.612960940022365]
Свободный член: -0.2209292751829534

63. Как преобразовать категориальные строковые данные в числовые индексы?

Уровень: L2

Входные данные:

data = [('cat',), ('dog',), ('mouse',), ('fish',), ('dog',), ('cat',), ('mouse',)]
df = spark.createDataFrame(data, ["animal"])
+------+
|animal|
+------+
|   cat|
|   dog|
| mouse|
|  fish|
|   dog|
|   cat|
| mouse|
+------+
Решение
from pyspark.ml.feature import StringIndexer

indexer = StringIndexer(inputCol='animal', outputCol='animalIndex')
indexed = indexer.fit(df).transform(df)
indexed.show()
+------+-----------+
|animal|animalIndex|
+------+-----------+
|   cat|        0.0|
|   dog|        1.0|
| mouse|        2.0|
|  fish|        3.0|
|   dog|        1.0|
|   cat|        0.0|
| mouse|        2.0|
+------+-----------+

64. Как вычислить корреляцию двух переменных в DataFrame?

Уровень: L1

Входные данные:

from pyspark.sql import Row

data = [
    Row(feature1=5, feature2=10, feature3=25),
    Row(feature1=6, feature2=15, feature3=35),
    Row(feature1=7, feature2=25, feature3=30),
    Row(feature1=8, feature2=20, feature3=60),
    Row(feature1=9, feature2=30, feature3=70),
]
df = spark.createDataFrame(data)
+--------+--------+--------+
|feature1|feature2|feature3|
+--------+--------+--------+
|       5|      10|      25|
|       6|      15|      35|
|       7|      25|      30|
|       8|      20|      60|
|       9|      30|      70|
+--------+--------+--------+
Решение
correlation = df.corr("feature1", "feature2")
print("Корреляция между feature1 и feature2:", correlation)
Корреляция между feature1 и feature2: 0.9

65. Как вычислить матрицу корреляции?

Уровень: L2

Входные данные:

from pyspark.sql import Row

data = [
    Row(feature1=5, feature2=10, feature3=25),
    Row(feature1=6, feature2=15, feature3=35),
    Row(feature1=7, feature2=25, feature3=30),
    Row(feature1=8, feature2=20, feature3=60),
    Row(feature1=9, feature2=30, feature3=70),
]
df = spark.createDataFrame(data)
+--------+--------+--------+
|feature1|feature2|feature3|
+--------+--------+--------+
|       5|      10|      25|
|       6|      15|      35|
|       7|      25|      30|
|       8|      20|      60|
|       9|      30|      70|
+--------+--------+--------+
Решение
from pyspark.ml.stat import Correlation
from pyspark.ml.feature import VectorAssembler

vector_assembler = VectorAssembler(
    inputCols=["feature1", "feature2", "feature3"], outputCol="features")
data_vector = vector_assembler.transform(df).select("features")

correlation_matrix = Correlation.corr(data_vector, "features").head()[0]
print(correlation_matrix)
DenseMatrix([[1.        , 0.9       , 0.91779992],
             [0.9       , 1.        , 0.67837385],
             [0.91779992, 0.67837385, 1.        ]])

66. Как вычислить VIF (коэффициент инфляции дисперсии)?

Уровень: L3

Входные данные:

from pyspark.sql import Row

data = [
    Row(feature1=5, feature2=10, feature3=25),
    Row(feature1=6, feature2=15, feature3=35),
    Row(feature1=7, feature2=25, feature3=30),
    Row(feature1=8, feature2=20, feature3=60),
    Row(feature1=9, feature2=30, feature3=70),
]
df = spark.createDataFrame(data)
+--------+--------+--------+
|feature1|feature2|feature3|
+--------+--------+--------+
|       5|      10|      25|
|       6|      15|      35|
|       7|      25|      30|
|       8|      20|      60|
|       9|      30|      70|
+--------+--------+--------+
Решение
from pyspark.ml.regression import LinearRegression
from pyspark.ml.feature import VectorAssembler

def calculate_vif(data, features):
    vif_dict = {}
    for feature in features:
        non_feature_cols = [col for col in features if col != feature]
        assembler = VectorAssembler(inputCols=non_feature_cols, outputCol="features")
        lr = LinearRegression(featuresCol='features', labelCol=feature)
        model = lr.fit(assembler.transform(data))
        vif = 1 / (1 - model.summary.r2)
        vif_dict[feature] = vif
    return vif_dict

features = ['feature1', 'feature2', 'feature3']
vif_values = calculate_vif(df, features)

for feature, vif in vif_values.items():
    print(f'VIF для {feature}: {vif:.2f}')
VIF для feature1: 66.21
VIF для feature2: 19.34
VIF для feature3: 23.30

67. Как выполнить критерий хи-квадрат (Chi-Square test)?

Уровень: L2

Входные данные:

data = [(1, 0, 0, 1, 1),
        (2, 0, 1, 0, 0),
        (3, 1, 0, 0, 0),
        (4, 0, 0, 1, 1),
        (5, 0, 1, 1, 0)]

df = spark.createDataFrame(data, ["id", "feature1", "feature2", "feature3", "label"])
+---+--------+--------+--------+-----+
| id|feature1|feature2|feature3|label|
+---+--------+--------+--------+-----+
|  1|       0|       0|       1|    1|
|  2|       0|       1|       0|    0|
|  3|       1|       0|       0|    0|
|  4|       0|       0|       1|    1|
|  5|       0|       1|       1|    0|
+---+--------+--------+--------+-----+
Решение
from pyspark.ml.feature import VectorAssembler
from pyspark.ml.stat import ChiSquareTest

assembler = VectorAssembler(inputCols=["feature1", "feature2", "feature3"], outputCol="features")
df = assembler.transform(df)

r = ChiSquareTest.test(df, "features", "label").head()
print("p-значения: " + str(r.pValues))
print("Степени свободы: " + str(r.degreesOfFreedom))
print("Статистики: " + str(r.statistics))
p-значения: [0.36131042852617856, 0.13603712811414348, 0.1360371281141436]
Степени свободы: [1, 1, 1]
Статистики: [0.8333333333333335, 2.2222222222222228, 2.2222222222222223]

68. Как вычислить стандартное отклонение?

Уровень: L1

Входные данные:

data = [("James",   "Sales",     3000),
        ("Michael", "Sales",     4600),
        ("Robert",  "Sales",     4100),
        ("Maria",   "Finance",   3000),
        ("James",   "Sales",     3000),
        ("Scott",   "Finance",   3300),
        ("Jen",     "Finance",   3900),
        ("Jeff",    "Marketing", 3000),
        ("Kumar",   "Marketing", 2000),
        ("Saif",    "Sales",     4100)]

df = spark.createDataFrame(data, ["Employee", "Department", "Salary"])
+--------+----------+------+
|Employee|Department|Salary|
+--------+----------+------+
|   James|     Sales|  3000|
| Michael|     Sales|  4600|
|  Robert|     Sales|  4100|
|   Maria|   Finance|  3000|
|   James|     Sales|  3000|
|   Scott|   Finance|  3300|
|     Jen|   Finance|  3900|
|    Jeff| Marketing|  3000|
|   Kumar| Marketing|  2000|
|    Saif|     Sales|  4100|
+--------+----------+------+
Решение
from pyspark.sql.functions import stddev

salary_stddev = df.select(stddev("Salary").alias("stddev"))
salary_stddev.show()
+-----------------+
|           stddev|
+-----------------+
|765.9416862050705|
+-----------------+

69. Как вычислить процент пропущенных значений в каждом столбце?

Уровень: L3

Входные данные:

data = [("John",  "Doe",   None),
        (None,    "Smith", "New York"),
        ("Mike",  "Smith", None),
        ("Anna",  "Smith", "Boston"),
        (None,    None,    None)]

df = spark.createDataFrame(data, ["FirstName", "LastName", "City"])
df.show()
+---------+--------+--------+
|FirstName|LastName|    City|
+---------+--------+--------+
|     John|     Doe|    null|
|     null|   Smith|New York|
|     Mike|   Smith|    null|
|     Anna|   Smith|  Boston|
|     null|    null|    null|
+---------+--------+--------+
Решение
total_rows = df.count()

for column in df.columns:
    null_values = df.filter(df[column].isNull()).count()
    missing_percentage = (null_values / total_rows) * 100
    print(f"Пропуски в {column}: {missing_percentage:.1f}%")
Пропуски в FirstName: 40.0%
Пропуски в LastName: 20.0%
Пропуски в City: 60.0%

70. Как получить имена всех DataFrame-объектов в текущем окружении?

Уровень: L2

import pyspark

dataframe_names = [
    name for name, obj in globals().items()
    if isinstance(obj, pyspark.sql.DataFrame)
]

for name in dataframe_names:
    print(name)