Сначала зафиксируйте grain данных, допущения, метрику и риск leakage, затем обсуждайте модель, инструмент или инфраструктуру.
Вопросы и ответы
10 подробных ответов
01Опишите архитектуру Spark: драйвер, экзекьюторы, cluster manager. Как задание разбивается на jobs, stages и tasks?
middle
Короткий ответ: Драйвер строит план и координирует выполнение, cluster manager (YARN, Kubernetes, Standalone) выделяет ресурсы, а экзекьюторы на воркерах выполняют задачи и хранят данные в памяти. Каждый экшен запускает job, который планировщик режет на стадии по границам шафла, а стадию — на задачи по числу партиций.
Подробно:
- Драйвер — процесс с вашим кодом и
SparkSession. Строит логический и физический план, DAG стадий, планирует задачи и собирает результаты. Падение драйвера убивает всё приложение. - Cluster manager — договаривается о ресурсах: сколько экзекьюторов, сколько ядер и памяти на каждый. Сам вычислениями не занимается.
- Экзекьюторы — JVM-процессы на воркерах. Выполняют задачи, кэшируют партиции, отдают данные через shuffle. Живут всё время работы приложения.
- Иерархия работы —
Job(на каждый экшен) →Stage(границы по шафлу) →Task(одна задача на одну партицию, минимальная единица параллелизма).
┌──────────────┐
│ Драйвер │ план + DAG + планировщик
│ SparkSession │
└──────┬───────┘
│ запрос ресурсов
┌──────▼───────┐
│Cluster Manager│ YARN / K8s / Standalone
└──────┬───────┘
┌──────────┼──────────┐
┌────▼────┐ ┌───▼────┐ ┌───▼────┐
│Экзекьютор│ │Экзекьютор│ │Экзекьютор│ задачи + кэш
└─────────┘ └────────┘ └────────┘
⚠️ Частая ошибка: путать ядра экзекьютора с задачами — параллелизм ограничен число экзекьюторов × ядер, и если партиций меньше, чем слотов, часть ядер простаивает.
02RDD против DataFrame/Dataset: почему на практике предпочитают DataFrame и когда всё-таки стоит спуститься до RDD?
middle
Короткий ответ: DataFrame — это структурированный API поверх RDD, за которым стоят оптимизатор Catalyst и движок Tungsten: они переписывают план и работают с колоночным бинарным представлением, поэтому DataFrame почти всегда быстрее и экономнее по памяти. К RDD спускаются, когда нужен полный контроль над низкоуровневой логикой или данные принципиально неструктурированы.
Подробно:
- Catalyst — оптимизатор запросов: pushdown предикатов и проекций, переупорядочивание джойнов, свёртка констант. Для RDD его нет — Spark выполняет ваш код как есть.
- Tungsten — off-heap бинарный формат и кодогенерация; меньше нагрузка на GC, чем у объектов JVM в RDD.
- Когда RDD оправдан — кастомная логика партиционирования, работа с совсем сырыми данными, тонкий контроль над низким уровнем или legacy-код.
| Критерий | RDD | DataFrame / Dataset |
|---|---|---|
| Оптимизатор | нет | Catalyst |
| Хранение | объекты JVM | Tungsten, колоночно |
| Схема | нет | есть, типизирована |
| Производительность | ниже | выше |
| Когда брать | сырые данные, контроль | 95% задач ETL/аналитики |
⚠️ Частая ошибка: без нужды писать пайплайн на RDD «ради гибкости» — вы отключаете Catalyst и Tungsten и почти всегда теряете в скорости на ровном месте.
03Что такое ленивые вычисления в Spark? Чем трансформации отличаются от экшенов и что такое узкие и широкие трансформации?
senior
Короткий ответ: Трансформации (map, filter, join) ленивы — они лишь дописывают шаги в DAG и ничего не считают, пока не вызван экшен (count, collect, write), который запускает job. Узкие трансформации не требуют перемешивания данных между партициями, широкие требуют шафла и потому образуют границу стадии.
Подробно:
- Ленивость — Spark копит план и оптимизирует его целиком через Catalyst перед запуском; лишние шаги вычёркиваются, предикаты проталкиваются к источнику.
- Трансформации против экшенов — трансформация возвращает новый DataFrame и ничего не запускает; экшен возвращает результат драйверу или пишет наружу и запускает вычисление.
- Узкие — каждая выходная партиция зависит от одной входной:
map,filter,union. Выполняются в конвейере внутри одной стадии, без сети. - Широкие — выходная партиция зависит от многих входных:
groupBy,join,distinct,repartition. Требуют шафла → новая стадия.
| Свойство | Узкая | Широкая |
|---|---|---|
| Зависимость партиций | 1 → 1 | много → 1 |
| Шафл | нет | да |
| Примеры | map, filter |
groupBy, join |
| Граница стадии | нет | да |
⚠️ Частая ошибка: думать, что df.filter(...) уже отфильтровал данные — до экшена не выполнено ничего, а повторные экшены пересчитывают всю цепочку заново, если её не закэшировать.
04Что такое shuffle в Spark, почему он дорогой, какие операции его вызывают и как его минимизировать?
senior
Короткий ответ: Шафл — это перераспределение данных между партициями и по сети так, чтобы записи с одним ключом оказались на одном экзекьюторе. Он дорогой из-за сериализации, дисковых записей на map-стороне и сетевой передачи, и это обычно узкое место джоба. Вызывают его широкие трансформации: groupBy, join, distinct, repartition.
Подробно:
- Почему дорого — map-задачи пишут отсортированные по ключу блоки на локальный диск (shuffle write), reduce-задачи тянут их по сети (shuffle read). Плюс сериализация и давление на GC.
- Что вызывает — любая широкая трансформация: агрегации, джойны,
distinct,orderBy, явныйrepartition. - Как минимизировать — broadcast-джойн для маленькой таблицы вместо sort-merge; фильтровать до джойна, а не после; избегать лишних
repartition; включить AQE, чтобы Spark сам сливал мелкие партиции после шафла.
Стадия map шафл Стадия reduce
[П0: a,b,a] ─┐ (запись на диск, ┌─► [ключ a: a,a,a]
[П1: b,c,a] ─┼─► передача по сети) ─┼─► [ключ b: b,b]
[П2: c,a,b] ─┘ пере-хэш по ключу └─► [ключ c: c,c]
⚠️ Частая ошибка: ставить repartition перед каждой операцией «для равномерности» — вы добавляете лишний полный шафл; часто хватает AQE или coalesce без перемешивания.
05Что такое перекос данных (data skew) в Spark, как он проявляется и как с ним бороться?
senior
Короткий ответ: Перекос — это неравномерное распределение данных по ключу, когда одна-две партиции получают несоразмерно много записей. Проявляется как одна «зависшая» задача: 199 задач стадии завершились, а 200-я тянется в разы дольше или падает по памяти. Лечится солением ключа, broadcast-джойном, AQE skew join и осмысленным репартиционированием.
Подробно:
- Как заметить — в Spark UI распределение длительности и shuffle read по задачам сильно неравномерно; max-задача на порядок дольше медианы.
- Salting — к горячему ключу добавляют случайный суффикс, чтобы разбить его на N партиций; на второй стороне джойна ключ размножают на те же N значений.
- AQE skew join — с Spark 3+ включённый
spark.sql.adaptive.skewJoin.enabledсам дробит перекошенные партиции во время выполнения. - Broadcast — если одна сторона маленькая, broadcast-джойн вовсе убирает шафл и перекос.
from pyspark.sql import functions as F
# солим горячий ключ на N корзин, чтобы разбить перекос
N = 16
facts = df.withColumn("salt", (F.rand() * N).cast("int"))
dims = dim.withColumn("salt", F.explode(F.array(
*[F.lit(i) for i in range(N)])))
result = facts.join(dims, ["key", "salt"], "inner")
⚠️ Частая ошибка: наращивать память экзекьюторов, чтобы «пережить» перекос — это лечит симптом ценой ресурсов; проблему решает соление или AQE, а не гигабайты heap.
06Как устроено партиционирование в Spark? Чем repartition отличается от coalesce и как выбрать число партиций?
middle
Короткий ответ: Партиция — минимальная единица параллелизма: одна задача обрабатывает одну партицию. repartition(n) делает полный шафл и может как увеличить, так и уменьшить число партиций с перемешиванием; coalesce(n) только уменьшает число партиций, сливая соседние без шафла, поэтому он дешевле, но может дать перекос.
Подробно:
- Число партиций — ориентир: партиция 128–256 МБ и хотя бы столько партиций, сколько всего слотов (
экзекьюторы × ядра), лучше кратно, чтобы не было хвостов. - repartition — полный шафл, равномерное распределение; берут, когда партиций слишком мало или нужно перераспределить по ключу (
repartition(col)). - coalesce — без шафла, только укрупняет; идеален перед записью, чтобы не плодить тысячи мелких файлов.
- Слишком много партиций — накладные расходы на планирование и мелкие файлы; слишком мало — недогруз параллелизма и spill.
| Критерий | repartition(n) |
coalesce(n) |
|---|---|---|
| Шафл | да, полный | нет |
| Направление | больше или меньше | только меньше |
| Равномерность | высокая | возможен перекос |
| Когда брать | поднять параллелизм | схлопнуть перед записью |
⚠️ Частая ошибка: делать repartition(1) перед записью, чтобы получить один файл — вы сгоняете все данные на один экзекьютор; для укрупнения используйте coalesce, а один файл собирайте только на действительно малых объёмах.
07Какие стратегии джойнов есть в Spark? Чем broadcast-джойн отличается от sort-merge и когда выигрывает broadcast?
middle
Короткий ответ: Sort-merge join шафлит обе таблицы по ключу, сортирует и сливает — универсальный вариант для двух больших датасетов. Broadcast (map-side) join рассылает маленькую таблицу целиком на все экзекьюторы, и джойн идёт локально без шафла большой стороны. Broadcast выигрывает, когда одна таблица помещается в память (по умолчанию до ~10 МБ, порог настраивается).
Подробно:
- Sort-merge — стратегия по умолчанию для больших-больших: два шафла, сортировка, слияние. Дорого, но масштабируется.
- Broadcast hash join — маленькая таблица собирается на драйвере и рассылается; на каждой партиции большой таблицы строится локальный хэш — шафла большой стороны нет.
- Когда broadcast — одна сторона мала (справочник, измерение); Spark с AQE часто сам переключается на broadcast, оценив размер во время выполнения.
from pyspark.sql import functions as F
# явно просим broadcast маленькой таблицы измерений
result = big_facts.join(
F.broadcast(small_dim),
on="dim_id",
how="inner",
)
# большая таблица не шафлится — джойн идёт на map-стороне
⚠️ Частая ошибка: броадкастить таблицу, которая не помещается в память драйвера/экзекьютора — это выбьет OOM; broadcast только для действительно маленькой стороны, порог — spark.sql.autoBroadcastJoinThreshold.
08Когда помогает cache()/persist() в Spark? Какие есть уровни хранения и чем плохо закэшировать не то?
middle
Короткий ответ: Кэширование окупается, когда один и тот же DataFrame используется несколько раз, — иначе Spark из-за ленивости пересчитывает всю цепочку на каждый экшен. cache() — это persist() с уровнем MEMORY_AND_DISK; persist(level) даёт выбор уровня. Кэш занимает память экзекьюторов, поэтому кэшировать «на всякий случай» вредно — вы вытесняете полезные данные и провоцируете spill.
Подробно:
- Когда помогает — итеративные алгоритмы, переиспользование промежуточного результата, ветвление плана из одной точки.
- Уровни хранения — компромисс между памятью, CPU и диском;
MEMORY_ONLYбыстрее, но при нехватке партиция теряется и пересчитывается. - Цена ошибки — лишний кэш съедает heap, растёт GC и spill; кэш нужно освобождать
unpersist(), когда он больше не нужен.
| Уровень | Где хранит | Комментарий |
|---|---|---|
| MEMORY_ONLY | RAM | быстро, но при нехватке — пересчёт |
| MEMORY_AND_DISK | RAM → диск | дефолт cache(), безопаснее |
| DISK_ONLY | диск | для больших, редко читаемых |
| *_SER | сериализованно | меньше памяти, больше CPU |
⚠️ Частая ошибка: кэшировать DataFrame, который читается ровно один раз — вы тратите память и время на материализацию, ничего не выигрывая; кэш имеет смысл только при повторном использовании.
09Как тюнить медленный Spark-джоб? Разберите размер партиций, spill, память, Adaptive Query Execution и чтение Spark UI.
senior
Короткий ответ: Начинают не с догадок, а со Spark UI: ищут стадию с самой длинной задачей, spill на диск и перекос в shuffle read. Затем правят размер партиций под 128–256 МБ, включают AQE, чтобы Spark сам подбирал число partition после шафла и переключал стратегию джойна, и балансируют память экзекьютора против spill.
Подробно:
- Читаем Spark UI — вкладка Stages: max/median duration задач (перекос), Spill (Memory/Disk) — признак нехватки памяти, Shuffle Read/Write — объём сети.
- Размер партиций — слишком крупные → spill и OOM; слишком мелкие → накладные расходы; цель 128–256 МБ.
- AQE —
spark.sql.adaptive.enabled=true: слияние мелких партиций, обработка skew join, авто-broadcast по фактическому размеру. - Память и spill — если много Disk Spill, поднять память/уменьшить партиции; следить за долей на shuffle.
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
# базовый параллелизм под ваши слоты
spark.conf.set("spark.sql.shuffle.partitions", "400")
⚠️ Частая ошибка: оставлять дефолтные spark.sql.shuffle.partitions=200 на больших данных — партиции становятся огромными и уходят в spill; подгоняйте число под объём и слоты (или доверьте это AQE).
10Объясните основы MapReduce (map, shuffle, reduce) и почему in-memory модель Spark быстрее на итеративных нагрузках.
middle
Короткий ответ: MapReduce обрабатывает данные в три фазы: map преобразует записи в пары ключ-значение, shuffle группирует их по ключу, reduce агрегирует. Классический Hadoop MapReduce материализует результат каждой фазы на диск (HDFS), а Spark держит промежуточные данные в памяти и строит DAG из множества стадий, поэтому на итеративных задачах он в разы быстрее.
Подробно:
- Map — параллельно по партициям превращает вход в пары ключ-значение (например,
(слово, 1)). - Shuffle — перераспределяет пары так, чтобы один ключ попал на один reducer; это сетевая и дисковая фаза.
- Reduce — сворачивает значения одного ключа (сумма, счётчик, конкатенация).
- Почему Spark быстрее — не пишет промежуток на HDFS между шагами, конвейеризует узкие трансформации, кэширует переиспользуемые данные и оптимизирует весь DAG через Catalyst.
Вход Map Shuffle Reduce
"a b a" ─► (a,1)(b,1)(a,1) ─► a:[1,1] b:[1] ─► a:2 b:1
"b c" ─► (b,1)(c,1) ─► c:[1] ─► c:1
⚠️ Частая ошибка: считать, что Spark «всегда в памяти и не пишет на диск» — при шафле и spill он тоже пишет на локальный диск; выигрыш в том, что между стадиями нет обязательной материализации в HDFS, а не в полном отказе от диска.
Источники
Источники и редакционная политика
Материалы RecallDeck сопоставлены с официальной документацией и открытыми публикациями компаний, когда первичный источник доступен. Мы не связаны с упомянутыми работодателями, не публикуем конфиденциальные задания и не продаём места в подборках. Формат найма может меняться — уточняйте его у рекрутера.
От чтения к воспроизведению
Отрепетируйте полный цикл интервью.
RecallDeck возвращает сложные темы по расписанию и помогает удерживать в памяти язык, SQL, архитектуру и поведенческие истории.