Блок 8 - Практические вопросы

60 вопросов в формате сертификационного экзамена Databricks: архитектура, кэширование, партиционирование, DataFrame API, join, UDF, I/O - с развёрнутыми ответами.

core exam

Вопросы составлены на основе официального практического экзамена Databricks Certified Associate Developer for Apache Spark 3.0 и дополнительных учебных материалов. Для каждого вопроса сначала попробуйте ответить самостоятельно, затем откройте ответ.


Блок 1 - Архитектура и модель выполнения

Вопрос № 1: Что из перечисленного описывает узкую трансформацию (narrow transformation)?

Варианты ответа:

  • A. Узкая трансформация - это операция, которая вычисляется на драйвере
  • B. Узкая трансформация - это операция, при которой данные не передаются между узлами кластера
  • C. Узкая трансформация всегда записывает данные на диск
  • D. Узкая трансформация создаёт новый DataFrame с другой схемой
  • E. Узкая трансформация - синоним Adaptive Query Execution
Ответ

Правильный ответ: B

При узкой трансформации каждая выходная партиция зависит ровно от одной входной партиции - данные не перемешиваются между executors. Примеры: filter(), map(), select(), withColumn(), coalesce(). Широкие трансформации (shuffle) - groupBy(), join(), repartition(), distinct() - требуют обмена данными между узлами.


Вопрос № 2: Какое из утверждений про stages в Spark корректно?

Варианты ответа:

  • A. Все задачи в stage выполняются последовательно
  • B. Stage всегда соответствует ровно одной JVM
  • C. Задачи внутри одного stage могут выполняться на нескольких машинах одновременно
  • D. Stage не может содержать более 200 задач
  • E. Stage создаётся только при чтении данных из внешнего хранилища
Ответ

Правильный ответ: C

Задачи (tasks) в stage независимы друг от друга и обрабатывают разные партиции - Spark распределяет их по executors и запускает параллельно. Stage завершается только когда завершены все его задачи. Граница между stages - это операция shuffle.


Вопрос № 3: Кто в архитектуре Spark назначает задачи executors?

Варианты ответа:

  • A. Cluster Manager
  • B. Worker Node
  • C. SparkContext
  • D. Driver
  • E. Executor сам запрашивает задачи
Ответ

Правильный ответ: D

Driver - главный процесс приложения. Он содержит DAG Scheduler и Task Scheduler, строит план выполнения и назначает задачи конкретным executors. Cluster Manager лишь выделяет ресурсы (память, CPU) и запускает executor-процессы - он не управляет выполнением задач.


Вопрос № 4: Чем отличается cluster mode от client mode?

Варианты ответа:

  • A. В cluster mode используется YARN, в client mode - Kubernetes
  • B. В cluster mode драйвер запускается на worker-узле; в client mode - на узле, с которого отправлено задание
  • C. В client mode доступно больше памяти для executor
  • D. В cluster mode нельзя использовать Python
  • E. Client mode и cluster mode - это синонимы
Ответ

Правильный ответ: B

В cluster mode драйвер запускается прямо на одном из worker-узлов кластера - удобно для продакшена, так как исключает сетевые задержки между драйвером и executors. В client mode драйвер работает на машине, с которой запущен spark-submit, - удобно для разработки и интерактивных сессий (ноутбуки).


Вопрос № 5: Какое из утверждений про shuffle корректно?

Варианты ответа:

  • A. Shuffle никогда не записывает данные на диск
  • B. Shuffle - это узкая трансформация
  • C. При shuffle Spark записывает данные на диск
  • D. Shuffle выполняется только при использовании sort()
  • E. Shuffle не влияет на производительность при наличии достаточного количества памяти
Ответ

Правильный ответ: C

В процессе shuffle Spark записывает промежуточные данные на диск (shuffle write), а затем читает их с диска другими executors (shuffle read). Это делает shuffle самой дорогой операцией. Широкие трансформации (groupBy, join, repartition, distinct) всегда вызывают shuffle.


Вопрос № 6: Что из перечисленного может инициировать Adaptive Query Execution (AQE)?

Варианты ответа:

  • A. Только трансформации
  • B. Только чтение данных из S3
  • C. Actions, но не трансформации
  • D. Только persist()
  • E. Любой вызов метода DataFrame
Ответ

Правильный ответ: C

AQE активируется при запуске action (например, count(), collect(), write()). Трансформации не запускают вычисления - Spark строит только логический план. Когда action инициирует выполнение, AQE может на лету менять стратегию join, объединять мелкие партиции после shuffle и оптимизировать skewed joins.


Вопрос № 7: Через что cluster manager получает информацию от драйвера?

Варианты ответа:

  • A. Через REST API
  • B. Через SparkContext
  • C. Через HDFS
  • D. Через Kafka
  • E. Cluster Manager не взаимодействует с Driver напрямую
Ответ

Правильный ответ: B

SparkContext - точка входа в Spark-приложение, создаётся на драйвере. Именно через SparkContext драйвер взаимодействует с cluster manager для запроса ресурсов, регистрации executors и координации выполнения.


Вопрос № 8: Какие режимы выполнения существуют в Spark?

Варианты ответа:

  • A. Client, Server
  • B. Interactive, Batch
  • C. Local, Remote
  • D. Client, Cluster, Local
  • E. Standalone, YARN, Kubernetes
Ответ

Правильный ответ: D

Три режима выполнения: local (всё на одной машине, для разработки), client (драйвер на клиентской машине, executors в кластере), cluster (драйвер и executors в кластере). Standalone, YARN, Mesos, Kubernetes - это режимы деплоя (cluster managers), не режимы выполнения.


Вопрос № 9: Может ли один Job в Spark охватывать несколько stages?

Варианты ответа:

  • A. Нет, один Job - всегда один Stage
  • B. Да, Job может охватывать несколько Stage
  • C. Только если используется persist()
  • D. Только при работе с RDD, но не с DataFrame
  • E. Только в cluster mode
Ответ

Правильный ответ: B

Job создаётся при вызове action и соответствует всему DAG от источника до результата. Этот DAG разбивается на несколько Stage по границам shuffle-операций. Например, Job с groupBy().agg().join() будет содержать минимум 3 Stage.


Вопрос № 10: Что из перечисленного инициирует выполнение (lazy evaluation) в Spark?

Варианты ответа:

  • A. Трансформации (transformations)
  • B. Определение схемы (StructType)
  • C. Actions
  • D. Импорт pyspark.sql.functions
  • E. Создание SparkSession
Ответ

Правильный ответ: C

Spark использует ленивые вычисления - трансформации только строят логический план, ничего не вычисляя. Вычисление запускается только когда вызывается action: count(), collect(), show(), write(), take() и т.д. Это позволяет Catalyst оптимизировать весь план перед выполнением.


Блок 2 - Кэширование и хранение

Вопрос № 11: Какой из вызовов немедленно удаляет закэшированный DataFrame df из памяти и диска?

Варианты ответа:

  • A. array_remove(df, '*')
  • B. df.unpersist()
  • C. del df
  • D. df.clearCache()
  • E. df.persist()
Ответ

Правильный ответ: B

df.unpersist() немедленно удаляет все закэшированные блоки DataFrame из памяти и диска.

  • del df - сообщает Python garbage collector о том, что объект можно удалить, но GC не запускается немедленно и кэш Spark не очищается.
  • df.clearCache() - такого метода у DataFrame нет.
  • df.persist() - наоборот, добавляет DataFrame в кэш.

Вопрос № 12: Какое из утверждений про уровень хранения MEMORY_AND_DISK некорректно?

Варианты ответа:

  • A. При нехватке памяти данные записываются на диск
  • B. MEMORY_AND_DISK реплицирует данные одновременно в память и на диск
  • C. MEMORY_AND_DISK - уровень по умолчанию для cache() на DataFrame
  • D. Для отказоустойчивого хранения следует использовать MEMORY_AND_DISK_2
  • E. MEMORY_AND_DISK не гарантирует наличие данных на всех узлах
Ответ

Правильный ответ: B

MEMORY_AND_DISK не реплицирует данные. Это означает: сначала партиции хранятся в памяти, а при нехватке места - спиллятся на диск того же executor. Репликация на 2 узлах обеспечивается только уровнем MEMORY_AND_DISK_2.


Вопрос № 13: Код должен кэшировать DataFrame transactionsDf в память executors и на диск при нехватке памяти, с поддержкой отказоустойчивости. Найдите ошибку:

transactionsDf.persist(StorageLevel.MEMORY_AND_DISK)

Варианты ответа:

  • A. persist() не принимает аргументов
  • B. Нужно использовать cache() вместо persist()
  • C. Уровень хранения не обеспечивает отказоустойчивость
  • D. MEMORY_AND_DISK не существует
  • E. Метод нужно вызвать дважды
Ответ

Правильный ответ: C

MEMORY_AND_DISK не обеспечивает репликацию - если executor упадёт, данные будут потеряны и Spark пересчитает партицию по lineage. Для fault-tolerant хранения нужен MEMORY_AND_DISK_2 (хранит копию данных на двух узлах):

transactionsDf.persist(StorageLevel.MEMORY_AND_DISK_2)

Вопрос № 14: Какой из вариантов сохраняет часть данных DataFrame itemsDf в памяти executors прямо сейчас?

Варианты ответа:

  • A. itemsDf.cache()
  • B. itemsDf.persist()
  • C. itemsDf.cache().count()
  • D. itemsDf.persist(StorageLevel.MEMORY_ONLY)
  • E. itemsDf.store().count()
Ответ

Правильный ответ: C

cache() - это ленивая трансформация: она лишь помечает DataFrame для кэширования, но ничего не вычисляет. Чтобы данные реально попали в память executors, нужен action. count() - action, который запускает вычисление и заполняет кэш. Вариант A только помечает, но не кэширует фактически.


Вопрос № 15: Какой из методов позволяет сохранить DataFrame только в памяти executor, если нужно явно настроить уровень хранения?

Варианты ответа:

  • A. df.clearCache()
  • B. df.storageLevel
  • C. df.cache()
  • D. df.persist()
  • E. StorageLevel.MEMORY_ONLY
Ответ

Правильный ответ: D

persist() принимает явный аргумент уровня хранения:

df.persist(StorageLevel.MEMORY_ONLY)

cache() не принимает аргументов и всегда использует MEMORY_AND_DISK. Если нужен другой уровень (например, MEMORY_ONLY, DISK_ONLY, OFF_HEAP) - только persist().


Вопрос № 16: Какой код корректно сохраняет DataFrame в памяти executor и при нехватке памяти - на диск?

Варианты ответа:

  • A. itemsDf.persist(StorageLevel.MEMORY_ONLY)
  • B. itemsDf.cache(StorageLevel.MEMORY_AND_DISK) - cache() не принимает аргументов
  • C. itemsDf.store()
  • D. itemsDf.cache()
  • E. itemsDf.write.option('destination', 'memory').save()
Ответ

Правильный ответ: D

cache() по умолчанию использует MEMORY_AND_DISK: сначала хранит в памяти, при нехватке - спиллирует на диск. Это именно то поведение, которое описано в вопросе. persist() с MEMORY_ONLY - неверно, так как при нехватке памяти данные будут пересчитываться, а не спиллироваться.


Вопрос № 17: Что произойдёт со значением аккумулятора, если action завершится неудачно, Spark перезапустит его и второй запуск пройдёт успешно?

Варианты ответа:

  • A. Значение аккумулятора удвоится (посчитается дважды)
  • B. Аккумулятор обнулится
  • C. В аккумулятор будет засчитана только успешная попытка
  • D. Аккумулятор будет содержать сумму обеих попыток
  • E. При сбое аккумулятор становится недействительным
Ответ

Правильный ответ: C

Spark гарантирует, что если action перезапускается после сбоя и завершается успешно, только успешная попытка будет учтена в аккумуляторе. Результаты неудачных попыток не добавляются. Это позволяет использовать аккумуляторы для надёжного подсчёта событий в распределённых вычислениях.


Блок 3 - Партиционирование и Shuffle

Вопрос № 18: Сколько партиций возвращает операция shuffle, если значение не задано явно?

Варианты ответа:

  • A. 10
  • B. 100
  • C. 300
  • D. 200
  • E. Зависит от размера данных
Ответ

Правильный ответ: D

Значение по умолчанию spark.sql.shuffle.partitions равно 200. Это количество партиций, которое создаётся после любой shuffle-операции (groupBy, join, repartition и т.д.). Для небольших датасетов 200 партиций - слишком много (много мелких задач); для больших - слишком мало.


Вопрос № 19: В коде ниже есть ошибка. Код должен задать количество shuffle-партиций равным 20. Найдите ошибку:

spark.conf.set(spark.sql.shuffle.partitions, 20)

Варианты ответа:

  • A. Значение 20 должно быть строкой "20"
  • B. Имя параметра должно быть передано как строка
  • C. Следует использовать spark.conf.get() вместо set()
  • D. Метод conf.set() принимает только один аргумент
  • E. Значение по умолчанию нельзя переопределять
Ответ

Правильный ответ: B

spark.conf.set() принимает строковое имя параметра. Без кавычек Python попытается разыменовать spark.sql.shuffle.partitions как цепочку атрибутов, что вызовет ошибку. Корректный вызов:

spark.conf.set("spark.sql.shuffle.partitions", 20)

Вопрос № 20: Какой параметр Spark определяет количество партиций при операциях join и aggregation?

Варианты ответа:

  • A. spark.sql.shuffle.partitions
  • B. spark.shuffle.partitions
  • C. spark.shuffle.io.maxRetries
  • D. spark.default.parallelism
  • E. spark.executor.cores
Ответ

Правильный ответ: A

spark.sql.shuffle.partitions управляет количеством post-shuffle партиций для DataFrame/SQL-операций. Значение по умолчанию: 200. spark.default.parallelism - аналогичный параметр, но для RDD API. spark.executor.cores - это количество CPU-ядер на executor, не количество партиций.


Вопрос № 21: Код должен вернуть 4-партиционный DataFrame из 8-партиционного storesDF без shuffle. Найдите ошибку:

storesDF.repartition(4)

Варианты ответа:

  • A. repartition работает только если DataFrame закэширован
  • B. repartition принимает только имя колонки, а не число
  • C. Из 8 партиций нельзя получить 4
  • D. repartition выполняет полный shuffle; нужен coalesce
  • E. repartition не гарантирует точное количество партиций
Ответ

Правильный ответ: D

repartition(n) всегда выполняет полный shuffle - данные перераспределяются по всем executors. Для уменьшения партиций без shuffle используется coalesce(n) - это узкая трансформация, которая объединяет существующие партиции без передачи данных по сети:

storesDF.coalesce(4)  # без shuffle

Вопрос № 22: Какой из вызовов всегда вернёт DataFrame из 12 партиций, если исходный DataFrame имеет 8 партиций?

Варианты ответа:

  • A. storesDF.coalesce(12) - coalesce может только уменьшать
  • B. storesDF.repartition() - нет аргумента
  • C. storesDF.repartition(12)
  • D. storesDF.coalesce() - нет аргумента
  • E. storesDF.coalesce(12, "storeId") - coalesce не принимает колонку
Ответ

Правильный ответ: C

repartition(n) может как увеличивать, так и уменьшать количество партиций путём полного shuffle. coalesce(n) может только уменьшать - она не может создать больше партиций, чем их есть сейчас.


Вопрос № 23: Является ли DataFrame.select() узкой или широкой трансформацией?

Варианты ответа:

  • A. Широкой, так как изменяет схему
  • B. Узкой
  • C. Зависит от количества выбранных колонок
  • D. Широкой, если выбирается более одной колонки
  • E. Это action, а не трансформация
Ответ

Правильный ответ: B

select() - узкая трансформация: каждая выходная партиция зависит ровно от одной входной партиции, данные не передаются между узлами. То же справедливо для filter(), withColumn(), drop(), alias(). Широкие трансформации требуют shuffle - например, groupBy(), join(), distinct().


Блок 4 - DataFrame API: Колонки и типы

Вопрос № 24: Какой код добавляет в DataFrame storesDF колонку numberOfManagers со значением константы 1 (целое число)?

Варианты ответа:

  • A. storesDF.withColumn("numberOfManagers", col(1))
  • B. storesDF.withColumn("numberOfManagers", 1)
  • C. storesDF.withColumn("numberOfManagers", lit(1))
  • D. storesDF.withColumn("numberOfManagers", lit("1"))
  • E. storesDF.withColumn("numberOfManagers", IntegerType(1))
Ответ

Правильный ответ: C

lit(value) оборачивает Python-литерал в Column-выражение. withColumn() ожидает второй аргумент типа Column, а не Python int или str. lit(1) создаёт integer-колонку; lit("1") создала бы string.

from pyspark.sql.functions import lit
storesDF.withColumn("numberOfManagers", lit(1))

Вопрос № 25: В коде ниже есть ошибка. Код должен разбить колонку storeCategory по символу _. Найдите ошибку:

storesDF.withColumn("storeValueCategory", col("storeCategory").split("_")[0])

Варианты ответа:

  • A. split() принимает строку-имя колонки, а не Column-объект
  • B. split() - это самостоятельная функция из pyspark.sql.functions, а не метод Column
  • C. Индексы [0] и [1] нужно передавать вторым аргументом в split()
  • D. Индексы должны быть 1 и 2, а не 0 и 1
  • E. withColumn() нельзя вызывать дважды подряд
Ответ

Правильный ответ: B

split() - самостоятельная функция из pyspark.sql.functions, а не метод класса Column. У объекта Column нет метода .split(). Корректный вариант:

from pyspark.sql.functions import split
storesDF.withColumn("storeValueCategory", split(col("storeCategory"), "_")[0])

Вопрос № 26: Какая операция создаёт по одной строке DataFrame для каждого элемента массива?

Варианты ответа:

  • A. extract()
  • B. split()
  • C. explode()
  • D. arrays_zip()
  • E. unpack()
Ответ

Правильный ответ: C

explode(col) превращает каждый элемент массива или ключ-значение из MapType в отдельную строку DataFrame. Если в строке было ["a", "b", "c"], после explode получится 3 строки. split() разбивает строку в массив; arrays_zip() объединяет массивы.


Вопрос № 27: Какой код корректно переводит колонку storeCategory в нижний регистр и заменяет её в DataFrame storesDF?

Варианты ответа:

  • A. storesDF.withColumn("storeCategory", lower(col("storeCategory")))
  • B. storesDF.withColumn("storeCategory", col("storeCategory").lower())
  • C. storesDF.withColumn("storeCategory", tolower(col("storeCategory")))
  • D. storesDF.withColumn("storeCategory", lower("storeCategory"))
  • E. storesDF.withColumn("storeCategory", lower(storeCategory))
Ответ

Правильный ответ: A

lower() - самостоятельная функция из pyspark.sql.functions. Она принимает Column-объект. У Column нет метода .lower() (вариант B). Функции tolower() в PySpark не существует (вариант C). Вариант D передаёт строку "storeCategory" - это работает как сокращение для col("storeCategory"), но вариант A явнее и корректнее.


Вопрос № 28: В коде ниже перепутан порядок аргументов в withColumnRenamed. Что является правилом для этого метода?

storesDF.withColumnRenamed("state", "division")

Варианты ответа:

  • A. Оба аргумента нужно обернуть в col()
  • B. withColumnRenamed() нельзя вызывать дважды; нужно использовать список
  • C. Старые колонки нужно явно удалять через drop()
  • D. Первый аргумент - старое имя колонки, второй - новое
  • E. withColumnRenamed нужно заменить на withColumn
Ответ

Правильный ответ: D

Сигнатура метода: withColumnRenamed(existingName, newName). В примере выше "state" - это existingName, а "division" - newName, то есть колонка state будет переименована в division. Аргументы не должны быть в col() - метод принимает строки.


Вопрос № 29: Как добавить в DataFrame itemsDf колонку itemNameParts - массив строк (максимум 4 элемента), разбитых по - или пробелу?

Варианты ответа:

  • A. itemsDf.withColumnRenamed("itemNameParts", split("itemName", "[\s\-]", 4))
  • B. itemsDf.withColumn("itemNameParts", split(col("itemName"), "[\s\-]", 4))
  • C. itemsDf.withColumn("itemNameParts", split(col("itemName"), "[\s\-]", 5))
  • D. itemsDf.withColumn("itemName", split(col("itemNameParts"), "[\s\-]", 4))
  • E. itemsDf.withColumn("itemNameParts", str_split(col("itemName"), "[\s\-]", 4))
Ответ

Правильный ответ: B

  • withColumn(newColName, expression) - добавляет новую колонку.
  • split(col, pattern, limit) - разбивает строку по паттерну; третий аргумент limit задаёт максимальное количество частей.
  • Для максимума 4 строк передаём 4.
  • str_split в PySpark не существует - только split.

Вопрос № 30: В коде ниже есть ошибка. Код должен переименовать колонку transactionId в transactionNumber. Найдите ошибку:

transactionsDf.withColumn("transactionNumber", "transactionId")

Варианты ответа:

  • A. Аргументы нужно поменять местами
  • B. Нужно поменять порядок аргументов и добавить copy()
  • C. Нужно добавить copy() в конец
  • D. Оба аргумента нужно обернуть в col(), а withColumn заменить на withColumnRenamed
  • E. withColumn нужно заменить на withColumnRenamed, а аргументы поменять местами
Ответ

Правильный ответ: E

Для переименования используется withColumnRenamed(existingName, newName). Метод withColumn(name, expression) предназначен для создания или замены колонки с Column-выражением, а не для переименования. Корректный код:

transactionsDf.withColumnRenamed("transactionId", "transactionNumber")

copy() в PySpark не нужен - все операции с DataFrame возвращают новый объект по умолчанию.


Блок 5 - Агрегации, сортировка, выборка

Вопрос № 31: Какой из вызовов удалит строки, где все колонки содержат null?

Варианты ответа:

  • A. storesDF.nadrop("all")
  • B. storesDF.na.drop("all", subset="sqft")
  • C. storesDF.dropna()
  • D. storesDF.na.drop()
  • E. storesDF.na.drop("all")
Ответ

Правильный ответ: E

  • na.drop("all") - удаляет строки, где все колонки null.
  • na.drop() без аргументов - удаляет строки, где хотя бы одна колонка null.
  • na.drop("all", subset=["col"]) - строки, где все колонки из subset равны null.
  • nadrop() и dropna() - в PySpark не существуют.

Вопрос № 32: Какой из вызовов завершится ошибкой при попытке удалить дублирующиеся строки?

Варианты ответа:

  • A. df.distinct()
  • B. df.drop_duplicates(subset=None)
  • C. df.drop_duplicates()
  • D. df.dropDuplicates()
  • E. df.drop_duplicates(subset="all")
Ответ

Правильный ответ: E

subset="all" - невалидный синтаксис. Параметр subset ожидает список имён колонок или None. Строка "all" не является зарезервированным значением и вызовет ошибку. Правильные варианты: distinct(), dropDuplicates(), drop_duplicates() - они эквивалентны.


Вопрос № 33: Какая функция не всегда возвращает точное количество уникальных значений в колонке?

Варианты ответа:

  • A. approx_count_distinct(col("division"))
  • B. countDistinct(col("division"))
  • C. df.select("division").dropDuplicates().count()
  • D. df.select("division").distinct().count()
  • E. Все перечисленные всегда точны
Ответ

Правильный ответ: A

approx_count_distinct() использует алгоритм HyperLogLog - он быстрее и дешевле, но возвращает приближённое (не точное) значение. Погрешность контролируется вторым параметром rsd (relative standard deviation). Все остальные варианты - точные.


Вопрос № 34: Какой код вычислит среднее значение колонки sqft и вернёт результат в колонке sqftMean?

Варианты ответа:

  • A. storesDF.agg(mean(col("sqft")).alias("sqftMean"))
  • B. storesDF.mean(col("sqft")).alias("sqftMean")
  • C. storesDF.withColumn("sqftMean", mean(col("sqft")))
  • D. storesDF.agg(average(col("sqft")).alias("sqftMean"))
  • E. storesDF.agg(mean("sqft", alias="sqftMean"))
Ответ

Правильный ответ: A

agg() принимает Column-выражения агрегации. mean() - функция из pyspark.sql.functions. Функции average() в PySpark не существует (вариант D - ошибка). Вариант C неверен: withColumn с агрегирующей функцией без groupBy не работает ожидаемым образом.


Вопрос № 35: Как получить количество строк в DataFrame storesDF?

Варианты ответа:

  • A. storesDF.withColumn("numberOfRows", count())
  • B. storesDF.countDistinct()
  • C. storesDF.agg(count("*"))
  • D. storesDF.count()
  • E. count(storesDF)
Ответ

Правильный ответ: D

DataFrame.count() - action, который возвращает количество строк как Python int. Вариант A создаёт колонку, а не считает строки. Вариант C вернёт однострочный DataFrame, а не число. У DataFrame нет метода countDistinct().


Вопрос № 36: Как получить сумму значений колонки sqft сгруппированную по колонке division?

Варианты ответа:

  • A. storesDF.groupBy.agg(sum(col("sqft")))
  • B. storesDF.groupBy("division").agg(sum())
  • C. storesDF.agg(groupBy("division").sum(col("sqft")))
  • D. storesDF.groupby("division").agg(sum(col("sqft")))
  • E. storesDF.groupBy("division").agg(sum(col("sqft")))
Ответ

Правильный ответ: E

groupBy() - метод, который нужно вызывать с аргументом и скобками. Вариант A (groupBy без скобок) вернёт GroupedData без вызова, что является ошибкой. Вариант D (groupby строчными) - в PySpark допустим как алиас, но вариант E является каноническим. Вариант B (sum() без аргумента) - ошибка.


Вопрос № 37: Какой код вернёт сводную статистику только для колонки sqft?

Варианты ответа:

  • A. storesDF.summary("mean")
  • B. storesDF.describe("sqft")
  • C. storesDF.summary(col("sqft"))
  • D. storesDF.describeColumn("sqft")
  • E. storesDF.summary()
Ответ

Правильный ответ: B

describe("sqft") возвращает DataFrame со статистиками count, mean, stddev, min, max только для колонки sqft. summary() принимает список статистик в виде строк ("mean", "stddev", "min", "max", "25%", "50%", "75%") - не колонки.


Вопрос № 38: Какие функции сортировки строк доступны в PySpark DataFrame?

Варианты ответа:

  • A. sort() и orderBy()
  • B. только orderBy()
  • C. только sort()
  • D. sortBy() и orderBy()
  • E. sort() и orderby() (строчными)
Ответ

Правильный ответ: A

DataFrame.sort() и DataFrame.orderBy() - это алиасы: оба сортируют DataFrame по указанным колонкам. orderby() строчными буквами - не существует. sortBy() - метод RDD, а не DataFrame.


Вопрос № 39: Код должен вернуть 15%-выборку строк без замены. Найдите ошибку:

storesDF.sample(True, fraction=0.15)

Варианты ответа:

  • A. Не указан аргумент seed
  • B. Не указан аргумент withReplacement
  • C. Для выборки без замены нужен метод sampleBy()
  • D. sample() не воспроизводим
  • E. Первый аргумент True задаёт выборку с заменой
Ответ

Правильный ответ: E

Сигнатура: sample(withReplacement, fraction, seed). Первый аргумент True устанавливает withReplacement=True - выборка с заменой. Для выборки без замены нужно передать False:

storesDF.sample(False, fraction=0.15)

Вопрос № 40: Как получить первые n строк DataFrame?

Варианты ответа:

  • A. df.n()
  • B. df.take(n)
  • C. df.head - атрибут, не вызов метода с n
  • D. df.show(n) - выводит, но возвращает None
  • E. df.collect(n) - collect не принимает аргументов
Ответ

Правильный ответ: B

take(n) и head(n) оба возвращают первые n строк как список объектов Row. show(n) - выводит в stdout и возвращает None. collect() не принимает аргументов и возвращает все строки - может вызвать OOM на больших DataFrame.


Вопрос № 41: Как получить значение поля sqft из первой строки DataFrame storesDF?

Варианты ответа:

  • A. storesDF.first.col("sqft")
  • B. storesDF.first.sqft
  • C. storesDF.first["sqft"]
  • D. storesDF.first()["sqft"]
  • E. storesDF.first().col("sqft")
Ответ

Правильный ответ: D

first() - это метод (нужны скобки), возвращает объект Row. К полям Row можно обращаться через ["field_name"] или через атрибут: row.sqft. Вариант B ошибочен: storesDF.first без скобок - это ссылка на метод, а не его вызов.


Блок 6 - Схема, UDF, SQL

Вопрос № 42: Какой вызов печатает схему DataFrame?

Варианты ответа:

  • A. print(df)
  • B. df.schema - возвращает объект StructType, не печатает
  • C. print(df.schema()) - schema это свойство, не метод
  • D. df.printSchema()
  • E. df.schema()
Ответ

Правильный ответ: D

printSchema() - action, который выводит схему в виде дерева в stdout. df.schema - свойство, возвращает объект StructType (его можно напечатать отдельно через print). schema() - ошибка синтаксиса, schema не является методом.


Вопрос № 43: В каком порядке нужно выполнить следующие строки для регистрации SQL-UDF ASSESS_PERFORMANCE и её применения?

1. spark.udf.register("ASSESS_PERFORMANCE", assessPerformance)
2. spark.sql("SELECT ASSESS_PERFORMANCE(customerSatisfaction) AS result FROM stores")
3. spark.udf.register(assessPerformance, "ASSESS_PERFORMANCE")
4. spark.sql("SELECT assessPerformance(customerSatisfaction) AS result FROM stores")

Варианты ответа:

  • A. 3, 4
  • B. 1, 2
  • C. 3, 2
  • D. 1, 4
  • E. 2, 1
Ответ

Правильный ответ: B

Шаги:

  1. Регистрация: spark.udf.register("ИМЯ", функция) - первый аргумент имя (строка), второй - функция.
  2. Вызов в SQL: используем зарегистрированное имя ASSESS_PERFORMANCE, а не оригинальное assessPerformance.

Вариант D ошибочен: assessPerformance недоступен в Spark SQL по имени Python-функции. Вариант C ошибочен: в строке 3 перепутан порядок аргументов.


Вопрос № 44: Как правильно создать Python UDF из функции assessPerformance возвращающей IntegerType, и применить её к DataFrame?

Варианты ответа:

  • A. assessPerformanceUDF = udf(assessPerformance, IntegerType()), затем df.withColumn("result", assessPerformanceUDF(col("customerSatisfaction")))
  • B. assessPerformanceUDF = spark.register.udf("ASSESS_PERFORMANCE", assessPerformance), затем df.withColumn(...)
  • C. assessPerformanceUDF = udf(assessPerformance, IntegerType) (без скобок)
  • D. assessPerformanceUDF = udf(assessPerformance), затем df.withColumn(...)
  • E. Напрямую: df.withColumn("result", assessPerformance(col("customerSatisfaction")))
Ответ

Правильный ответ: A

udf(function, returnType()) - тип возвращаемого значения должен быть инстанциирован с (). IntegerType без скобок - это класс, а не экземпляр, что вызовет ошибку. Вариант E не работает: нельзя применять Python-функцию напрямую к Column-объектам.


Вопрос № 45: Какая функция выполняет SQL-запрос на зарегистрированной таблице?

Варианты ответа:

  • A. spark.query()
  • B. DataFrame.sql()
  • C. spark.sql()
  • D. DataFrame.createOrReplaceTempView()
  • E. DataFrame.createTempView()
Ответ

Правильный ответ: C

spark.sql("SELECT ...") выполняет SQL-запрос и возвращает DataFrame. createOrReplaceTempView() и createTempView() - это регистрация временного представления, не выполнение запроса.


Вопрос № 46: Как правильно определить схему для колонки с массивом строк?

itemsDfSchema = StructType([
    StructField('itemId', IntegerType()),
    StructField('attributes', ???),
    StructField('supplier', StringType())
])

Варианты ответа:

  • A. StructField('attributes', StringType())
  • B. StructField('attributes', ArrayType(StringType))
  • C. StructField('attributes', array<string>)
  • D. StructField('attributes', ArrayType(StringType()))
  • E. StructField('attributes', List(StringType()))
Ответ

Правильный ответ: D

ArrayType(StringType()) - оба должны быть инстанциированы с (). ArrayType(StringType) без скобок у StringType не сработает. В DDL-строке это выглядит как array<string>, но не как аргумент StructField. List() в PySpark не существует.


Вопрос № 47: Как добавить колонку predErrorSqrt - квадратный корень из predError - в DataFrame transactionsDf?

Варианты ответа:

  • A. transactionsDf.withColumn("predErrorSqrt", sqrt(predError))
  • B. transactionsDf.select(sqrt(predError))
  • C. transactionsDf.withColumn("predErrorSqrt", col("predError").sqrt())
  • D. transactionsDf.withColumn("predErrorSqrt", sqrt(col("predError")))
  • E. transactionsDf.select(sqrt("predError"))
Ответ

Правильный ответ: D

  • sqrt() - самостоятельная функция из pyspark.sql.functions.
  • Колонку нужно передать как Column: col("predError").
  • Вариант A: predError без col() - Python попытается разыменовать переменную.
  • Вариант C: у Column нет метода .sqrt().
  • Варианты B и E возвращают DataFrame только с колонкой sqrt, без остальных колонок.

Блок 7 - Join и объединение

Вопрос № 48: Какой тип join использует DataFrame.join() по умолчанию?

Варианты ответа:

  • A. Left outer
  • B. Inner
  • C. Full outer
  • D. Cross
  • E. Left semi
Ответ

Правильный ответ: B

По умолчанию df.join(other, on) выполняет inner join - возвращает только строки, которые есть в обоих DataFrame. Явно указывать how="inner" необязательно, но рекомендуется для читаемости.


Вопрос № 49: Как выполнить outer join между storesDF и employeesDF по колонке storeId?

Варианты ответа:

  • A. storesDF.join(employeesDF, "storeId", "outer")
  • B. storesDF.join(employeesDF, "storeId")
  • C. storesDF.join(employeesDF, "outer", col("storeId"))
  • D. storesDF.join(employeesDF, "outer", storesDF.storeId == employeesDF.storeId)
  • E. storesDF.merge(employeesDF, "outer", col("storeId"))
Ответ

Правильный ответ: A

Сигнатура: join(other, on, how). Аргументы по порядку: другой DataFrame, условие/колонка join, тип join. Вариант B - inner join (нет how). Варианты C и D - перепутан порядок аргументов. merge() в PySpark не существует.


Вопрос № 50: В коде ниже есть ошибка. Код должен выполнить inner join storesDF и employeesDF по колонкам storeId и employeeId соответственно. Найдите ошибку:

storesDF.join(employeesDF, [col("storeId"), col("employeeId")])

Варианты ответа:

  • A. join() - самостоятельная функция, а не метод DataFrame
  • B. Нужен третий аргумент для типа join
  • C. col("storeId") и col("employeeId") нужно сравнить через ==
  • D. Метода DataFrame.join() нет, используйте DataFrame.merge()
  • E. Ссылки на колонки не нужно оборачивать в col() - нужно убрать col()
Ответ

Правильный ответ: C

При join по колонкам с разными именами условие задаётся через равенство Column-объектов:

storesDF.join(employeesDF, storesDF.storeId == employeesDF.employeeId)

Список [col("storeId"), col("employeeId")] интерпретируется как join по двум одноимённым колонкам, что здесь неприменимо.


Вопрос № 51: Какой параметр Spark управляет автоматическим broadcast join без явного вызова broadcast()?

Варианты ответа:

  • A. spark.sql.autoBroadcastJoinThreshold
  • B. spark.sql.broadcastTimeout
  • C. spark.broadcast.blockSize
  • D. spark.broadcast.compress
  • E. spark.executor.memoryOverhead
Ответ

Правильный ответ: A

spark.sql.autoBroadcastJoinThreshold задаёт максимальный размер DataFrame (в байтах), при котором Spark автоматически выберет broadcast join. Значение по умолчанию: 10 МБ (10485760 байт). Для установки 20 МБ: spark.conf.set("spark.sql.autoBroadcastJoinThreshold", 20 * 1024 * 1024).


Вопрос № 52: В коде ниже есть ошибка. Код должен настроить автоматический broadcast join для DataFrame размером до 20 МБ:

spark.conf.set("spark.sql.autoBroadcastJoinThreshold", 20)

Варианты ответа:

  • A. Имя параметра написано неверно
  • B. Значение 20 задаёт 20 байт, а не 20 МБ - параметр ожидает байты
  • C. Параметр нельзя менять через conf.set()
  • D. Значение должно быть строкой: "20mb"
  • E. Ошибки нет
Ответ

Правильный ответ: B

Параметр autoBroadcastJoinThreshold ожидает значение в байтах. Передача 20 устанавливает порог в 20 байт - намного меньше дефолтных 10 МБ. Корректный вариант:

spark.conf.set("spark.sql.autoBroadcastJoinThreshold", 20 * 1024 * 1024)  # 20 МБ

Вопрос № 53: Как выполнить cross join между storesDF и employeesDF?

Варианты ответа:

  • A. storesDF.crossJoin(employeesDF, "storeId")
  • B. storesDF.join(employeesDF, "cross")
  • C. storesDF.join(employeesDF, how="cross")
  • D. storesDF.join(employeesDF, col("storeId"), "cross")
  • E. storesDF.crossJoin(employeesDF)
Ответ

Правильный ответ: E

crossJoin(other) выполняет декартово произведение - каждая строка из первого DataFrame соединяется с каждой строкой второго. Метод не принимает аргумент с колонкой. Результат содержит m × n строк. Вариант C также корректен синтаксически, но crossJoin() - идиоматический способ.


Вопрос № 54: В чём разница между DataFrame.union() и DataFrame.unionByName()?

Варианты ответа:

  • A. union() работает только с одинаковым числом строк
  • B. union() объединяет по позиции колонок; unionByName() - по имени колонок
  • C. unionByName() удаляет дубликаты, union() - нет
  • D. Это синонимы
  • E. union() работает только с RDD
Ответ

Правильный ответ: B

union() совмещает DataFrame по позиции - первая колонка к первой, вторая ко второй, вне зависимости от имён. unionByName() совмещает по именам - удобно, когда порядок колонок в двух DataFrame разный. Оба метода не удаляют дубликаты.


Блок 8 - Чтение и запись данных

Вопрос № 55: Как записать DataFrame storesDF в Parquet по пути filePath?

Варианты ответа:

  • A. storesDF.write.option("parquet").path(filePath)
  • B. storesDF.write.path(filePath)
  • C. storesDF.write().parquet(filePath)
  • D. storesDF.write(filePath)
  • E. storesDF.write.parquet(filePath)
Ответ

Правильный ответ: E

DataFrame.write - свойство (не метод, без скобок), возвращает DataFrameWriter. .parquet(path) записывает данные в формате Parquet. Вариант C ошибочен: write - это свойство, а не вызываемый объект.


Вопрос № 56: В коде ниже есть ошибка. Код должен записать storesDF в Parquet, партиционированный по division. Найдите ошибку:

storesDF.write.repartition("division").parquet(filePath)

Варианты ответа:

  • A. division нужно обернуть в col()
  • B. Для DataFrameWriter нет .parquet() - нужен .save()
  • C. У DataFrameWriter нет .repartition() - нужен .partitionBy()
  • D. После write нужно поставить скобки ()
  • E. Нужно указать .mode() для перезаписи
Ответ

Правильный ответ: C

DataFrameWriter не имеет метода repartition(). Для партиционирования при записи (создание директорий по значениям колонки) используется partitionBy():

storesDF.write.partitionBy("division").parquet(filePath)

DataFrame.repartition() - это трансформация для перераспределения данных в памяти, а не инструмент записи.


Вопрос № 57: Как прочитать Parquet-файл по пути filePath в DataFrame?

Варианты ответа:

  • A. spark.read().parquet(filePath)
  • B. spark.read().path(filePath, source="parquet")
  • C. spark.read.path(filePath, source="parquet")
  • D. spark.read.parquet(filePath)
  • E. spark.read().path(filePath)
Ответ

Правильный ответ: D

spark.read - свойство (без скобок), возвращает DataFrameReader. .parquet(path) читает Parquet. Эквивалентный вариант через универсальный метод: spark.read.format("parquet").load(filePath).


Вопрос № 58: Как прочитать JSON-файл по filePath с заранее определённой схемой schema?

Варианты ответа:

  • A. spark.read().schema(schema).format(json).load(filePath)
  • B. spark.read().schema(schema).format("json").load(filePath)
  • C. spark.read.schema("schema").format("json").load(filePath)
  • D. spark.read.schema(schema).format("json").load(filePath)
  • E. spark.read.json(filePath, schema=schema)
Ответ

Правильный ответ: D

  • spark.read - свойство, не вызов (без скобок).
  • schema(schemaObject) - принимает объект StructType, не строку "schema".
  • format("json") - строка с именем формата.

Корректный код:

spark.read.schema(schema).format("json").load(filePath)


Вопрос № 59: Как прочитать Parquet-файл с объединением схем двух партиций, у которых разные наборы колонок?

Варианты ответа:

  • A. spark.read.parquet(filePath) - по умолчанию схемы объединяются
  • B. spark.read.option("mergeSchema", True).parquet(filePath)
  • C. spark.read.option("inferSchema", True).parquet(filePath)
  • D. spark.read.schema(schema1 + schema2).parquet(filePath)
  • E. Spark не поддерживает объединение схем
Ответ

Правильный ответ: B

Опция mergeSchema = true включает объединение схем при чтении Parquet: колонки из всех партиций объединяются в единую схему (колонки, которых нет в отдельных файлах, получают null). Ограничение: если одноимённая колонка имеет разные типы в разных партициях - mergeSchema завершится ошибкой.


Вопрос № 60: В коде ниже есть ошибка. Код должен записать Parquet, партиционированный по storeId. Найдите ошибку:

transactionsDf.write.partitionOn("storeId").parquet(filePath)

Варианты ответа:

  • A. storeId нужно передавать как col("storeId")
  • B. Режим записи нужно явно указать через mode()
  • C. parquet() не принимает путь - нужен save(filePath)
  • D. Метода partitionOn() не существует - нужен partitionBy()
  • E. Ошибки нет
Ответ

Правильный ответ: D

У DataFrameWriter нет метода partitionOn(). Правильный метод - partitionBy():

transactionsDf.write.partitionBy("storeId").parquet(filePath)

Блок 9 - Дата и время

Вопрос № 61: В каком порядке нужно выполнить следующие строки, чтобы получить колонку openDateString в формате "Sunday, Dec 4, 2008 1:05 PM" из UNIX timestamp?

1. storesDF.withColumn("openDateString", from_unixtime(col("openDate"), simpleDateFormat))
2. simpleDateFormat = "EEEE, MMM d, yyyy h:mm a"
3. storesDF.withColumn("openDateString", from_unixtime(col("openDate"), SimpleDateFormat()))
4. storesDF.withColumn("openDateString", date_format(col("openDate"), simpleDateFormat))
5. simpleDateFormat = "wd, MMM d, yyyy h:mm a"

Варианты ответа:

  • A. 2, 3
  • B. 2, 1
  • C. 5, 4
  • D. 2, 4
  • E. 5, 1
Ответ

Правильный ответ: B

Шаги:

  1. Определить формат: simpleDateFormat = "EEEE, MMM d, yyyy h:mm a" - EEEE даёт полное название дня недели.
  2. Применить: from_unixtime(col("openDate"), simpleDateFormat) - конвертирует UNIX epoch (секунды) в строку по паттерну.

SimpleDateFormat() - Java-класс, не Python-функция. date_format() работает с Timestamp, а не с integer-столбцом epoch.


Вопрос № 62: Как извлечь номер месяца из колонки openDate типа integer (UNIX epoch секунды)?

Варианты ответа:

  • A. df.withColumn("month", getMonth(col("openDate")))
  • B. df.withColumn("openTimestamp", col("openDate").cast("Timestamp")).withColumn("month", month(col("openTimestamp")))
  • C. df.withColumn("openDateFormat", col("openDate").cast("Date")).withColumn("month", month(col("openDateFormat")))
  • D. df.withColumn("month", substr(col("openDate"), 4, 2))
  • E. df.withColumn("month", month(col("openDate")))
Ответ

Правильный ответ: B

Функция month() работает с типом Timestamp, но не с integer. Необходимо сначала привести integer к Timestamp через .cast("Timestamp"), а затем применить month(). Приведение к Date (вариант C) тоже работает, но month() ожидает именно Timestamp. Вариант E - ошибка: integer не поддерживается month() напрямую.


Вопрос № 63: Какая функция конвертирует UNIX timestamp (integer) в строку? unix_timestamp() или to_unixtime()?

Ответ

Правильный ответ: unix_timestamp()

Функция to_unixtime() в PySpark не существует. Корректные функции для работы с UNIX timestamp:

  • unix_timestamp(timestamp, format) - конвертирует строку в UNIX epoch (секунды)
  • from_unixtime(unixTime, format) - конвертирует UNIX epoch в строку

Не перепутайте: from_unixtime() существует, а to_unixtime() - нет.


Блок 10 - Разное

Вопрос № 64: Как создать однострочный DataFrame из списка целых чисел years = [2018, 2019, 2020]?

Варианты ответа:

  • A. spark.createDataFrame([years], IntegerType())
  • B. spark.createDataFrame(years, IntegerType())
  • C. spark.DataFrame(years, IntegerType())
  • D. spark.createDataFrame(years)
  • E. spark.createDataFrame(years, IntegerType)
Ответ

Правильный ответ: B

createDataFrame(data, schema) - для одноколоночного DataFrame с числами можно передать IntegerType() как схему. Тип должен быть инстанциирован (IntegerType(), а не IntegerType). Вариант A - лишние скобки вокруг years создадут вложенный список. Вариант D вызовет ошибку без схемы.


Вопрос № 65: Доступен ли Dataset API в PySpark?

Варианты ответа:

  • A. Да, полностью доступен
  • B. Да, но только в PySpark 3.x
  • C. Нет, Dataset API доступен только в Scala и Java
  • D. Да, через отдельный модуль pyspark.dataset
  • E. Нет, Dataset API был удалён в Spark 3.0
Ответ

Правильный ответ: C

Dataset API предоставляет типизированные, compile-time безопасные трансформации. Это возможно только в статически типизированных языках - Scala и Java. В Python нет compile-time type checking, поэтому Dataset API недоступен. В PySpark эту роль частично играет использование Pandas UDF и схем StructType.


Вопрос № 66: Чем Spark DataFrame отличается от Pandas DataFrame?

Варианты ответа:

  • A. Они полностью эквивалентны
  • B. Spark DataFrame не поддерживает SQL
  • C. Spark DataFrame - распределённый, Pandas - нет; Spark DataFrame иммутабелен
  • D. Pandas DataFrame быстрее Spark при работе с большими данными
  • E. Spark DataFrame работает только с числовыми данными
Ответ

Правильный ответ: C

Ключевые отличия:

  • Распределённость: Spark DataFrame хранится на кластере в партициях; Pandas - только в памяти одной машины.
  • Иммутабельность: Spark DataFrame неизменяемы - каждая трансформация создаёт новый объект.
  • Ленивые вычисления: Spark откладывает выполнение до action.
  • Масштаб: Pandas ограничен RAM одной машины; Spark - масштабируется горизонтально.

Вопрос № 67: Какая операция из перечисленных вызывает наибольший сетевой трафик?

Варианты ответа:

  • A. df.count()
  • B. df.show(10)
  • C. df.write.parquet(path)
  • D. df.collect()
  • E. df.cache().count()
Ответ

Правильный ответ: D

collect() переносит все данные DataFrame с executors на драйвер. На большом Dataset это может означать гигабайты данных по сети и риск OOM на драйвере. count() передаёт только агрегированное число. show(10) - только 10 строк. write сохраняет данные напрямую из executors в хранилище.


Вопрос № 68: Какое из утверждений про аккумуляторы корректно?

Варианты ответа:

  • A. Аккумуляторы могут читать как driver, так и executors
  • B. Значение аккумулятора можно читать только на драйвере
  • C. Аккумуляторы обновляются только при вызове collect()
  • D. Аккумуляторы изменяемы - executors могут как читать, так и писать
  • E. Аккумуляторы - это broadcast-переменные в обратном направлении
Ответ

Правильный ответ: B

Аккумуляторы работают по принципу write-only на executors, read-only на driver: executors могут только прибавлять к значению, но не читать его. Только драйвер может прочитать итоговое значение через .value. Это делает их безопасными для агрегации без синхронизации.


Вопрос № 69: Как правильно зарегистрировать UDF для использования в Spark SQL?

Варианты ответа:

  • A. spark.udf.register(to_limit, "LIMIT_FCN")
  • B. spark.udf.register("LIMIT_FCN", to_limit)
  • C. udf("LIMIT_FCN", to_limit)
  • D. spark.sql.register("LIMIT_FCN", to_limit)
  • E. spark.register.udf(to_limit, "LIMIT_FCN")
Ответ

Правильный ответ: B

spark.udf.register(name, function) - первый аргумент строка с именем, второй - Python-функция. После регистрации UDF вызывается в SQL по зарегистрированному имени:

spark.udf.register("LIMIT_FCN", to_limit)
spark.sql("SELECT LIMIT_FCN(predError) AS result FROM transactionsDf")

В SQL нельзя вызывать функцию по её Python-имени - только по зарегистрированному.


Вопрос № 70: Какие методы переименовывают колонку, принимают ровно два строковых аргумента и не требуют оборачивания в col()?

Ответ

Правильный ответ: withColumnRenamed(existingName, newName)

  • Принимает ровно два аргумента-строки: старое имя и новое.
  • Аргументы не оборачиваются в col().
  • Для переименования нескольких колонок - цепочка вызовов:
df.withColumnRenamed("attributes", "feature0") \
  .withColumnRenamed("supplier", "feature1")

Частая ошибка: withColumnRenamed(col("oldName"), col("newName")) - ошибка, col() здесь не нужен. Другая ошибка: передать 4 аргумента для переименования двух колонок - метод принимает только 2.