Блок 8 - Практические вопросы
60 вопросов в формате сертификационного экзамена Databricks: архитектура, кэширование, партиционирование, DataFrame API, join, UDF, I/O - с развёрнутыми ответами.
Вопросы составлены на основе официального практического экзамена 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
Шаги:
- Регистрация:
spark.udf.register("ИМЯ", функция)- первый аргумент имя (строка), второй - функция. - Вызов в 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
Шаги:
- Определить формат:
simpleDateFormat = "EEEE, MMM d, yyyy h:mm a"-EEEEдаёт полное название дня недели. - Применить:
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.