SCALA / FLINK / ICEBERG = STUDY-GAPSCALA / FLINK / ICEBERG = ПРОБЕЛ ДЛЯ ИЗУЧЕНИЯ

Apache Spark — DE Interview PrepПодготовка к DE-интервью

Fast-revision cheatsheet for Data Engineering interviews. Tailored for Senior/Lead DE. Dense, skimmable, with interviewer-probe callouts, comparison tables, ASCII diagrams, and self-quiz flip cards. Anchor Spark hands-on to Front Tier (Scala/Spark on Hadoop); anchor scale to Meta (3.6B events/day). Шпаргалка для быстрого повторения перед Data Engineering интервью. Заточена под Senior/Lead DE. Плотная, удобная для беглого чтения, с указаниями на вопросы интервьюера, сравнительными таблицами, ASCII-диаграммами и карточками самопроверки. Практический якорь Spark — Front Tier (Scala/Spark на Hadoop); якорь масштаба — Meta (3,6 млрд событий/день).

1 · Fundamentals & Architecture1 · Основы и архитектура

The mental model in one screenМентальная модель на одном экране

SparkSession (Driver JVM) │ builds logical+physical plan, owns DAG/Task scheduler ▼ Cluster Manager (YARN / Kubernetes / Standalone) │ allocates executors ▼ ┌───────────────┐ ┌───────────────┐ ┌───────────────┐ │ Executor 1 │ │ Executor 2 │ │ Executor N │ │ tasks+cache │ │ tasks+cache │ │ tasks+cache │ │ shuffle disk │ │ shuffle disk │ │ shuffle disk │ └───────────────┘ └───────────────┘ └───────────────┘ Job ──split-at-shuffle──► Stages ──one-per-partition──► Tasks
SparkSession (Driver JVM) │ строит логический+физический план, владеет DAG/Task scheduler ▼ Cluster Manager (YARN / Kubernetes / Standalone) │ выделяет executors ▼ ┌───────────────┐ ┌───────────────┐ ┌───────────────┐ │ Executor 1 │ │ Executor 2 │ │ Executor N │ │ задачи+кеш │ │ задачи+кеш │ │ задачи+кеш │ │ shuffle disk │ │ shuffle disk │ │ shuffle disk │ └───────────────┘ └───────────────┘ └───────────────┘ Job ──разделение на shuffle──► Stages ──одна на партицию──► Tasks
TermWhat it is
DriverRuns main, builds plan, schedules, collects results. One per app. Don't collect() big data here.
ExecutorJVM on a worker; runs tasks, holds cache + shuffle files. Many per app.
JobTriggered by one action.
StageSet of tasks with no shuffle between them. Boundary = wide dependency.
TaskSmallest unit; one task per partition of a stage.
ТерминЧто это
DriverИсполняет main, строит план, планирует, собирает результаты. Один на приложение. Не делай collect() больших данных сюда.
ExecutorJVM на воркере; выполняет задачи, хранит кеш + shuffle-файлы. Много на приложение.
JobЗапускается одним action.
StageНабор задач без shuffle между ними. Граница = широкая зависимость.
TaskНаименьшая единица; одна задача на партицию стадии.

2 · RDD vs DataFrame vs Dataset2 · RDD vs DataFrame vs Dataset

RDDDataFrameDataset
SchemaNoYes (Row)Yes (typed)
Type safetyCompile-timeRuntime onlyCompile-time
Catalyst/TungstenNoYesYes
LanguagesAllAllScala/Java only
Use whenUnstructured, custom partitioner, full controlDefault ETL/analyticsTyped pipelines, testability
RDDDataFrameDataset
СхемаНетДа (Row)Да (типизированная)
ТипобезопасностьCompile-timeТолько runtimeCompile-time
Catalyst/TungstenНетДаДа
ЯзыкиВсеВсеТолько Scala/Java
Когда использоватьНеструктурированные данные, кастомный партиционер, полный контрольСтандартный ETL/аналитикаТипизированные пайплайны, тестируемость

One-liner: "Prefer DataFrame/Dataset so Catalyst can optimize; drop to RDD only when you need control the optimizer can't give you."

В одну строку: «Предпочитай DataFrame/Dataset, чтобы Catalyst мог оптимизировать; спускайся к RDD только когда нужен контроль, который оптимизатор дать не может.»

Anchor Sell the typed Dataset (Scala) as a craftsmanship/testing win over untyped PySpark: compile-time safety, refactorable, fewer runtime surprises.
Якорь Продавай типизированный Dataset (Scala) как преимущество качества/тестирования над нетипизированным PySpark: безопасность на этапе компиляции, рефакторимость, меньше runtime-сюрпризов.

3 · Transformations, Actions, Lazy Eval & the DAG3 · Трансформации, actions, ленивое выполнение и DAG

Transformations (lazy)Трансформации (ленивые)

Build the DAG, return a new RDD/DF, execute nothing.

Строят DAG, возвращают новый RDD/DF, ничего не выполняют.

map filter select join groupBy withColumn repartition

Actions (eager)Actions (нетерпеливые)

Trigger execution; return to driver or write out.

Запускают выполнение; возвращают драйверу или пишут наружу.

count collect take show write foreach reduce first

What happens on an action (the end-to-end walk)Что происходит при action (полный путь)

Lazy transforms ─► Unresolved Logical Plan ─► Analyzed (resolve vs catalog) ─► Optimized Logical Plan (Catalyst rules: pushdown, pruning, fold...) ─► Physical Plans ─► Cost-based pick ─► WholeStageCodegen ─► DAGScheduler splits into STAGES at shuffle boundaries ─► TaskScheduler ships TASKS to executors (1 / partition) ─► shuffle: map writes files to local disk ► reduce fetches over network ─► result to driver (collect) / sink (write)
Ленивые трансформации ─► Неразрешённый логический план ─► Проанализированный (резолв по каталогу) ─► Оптимизированный логический план (правила Catalyst: pushdown, pruning, fold...) ─► Физические планы ─► Выбор по стоимости ─► WholeStageCodegen ─► DAGScheduler делит на СТАДИИ на границах shuffle ─► TaskScheduler отправляет ЗАДАЧИ на executors (1 / партицию) ─► shuffle: map пишет файлы на локальный диск ► reduce забирает по сети ─► результат драйверу (collect) / в приёмник (write)
Gotcha An action re-runs the entire lineage every time unless you cache(). count() then write() on an uncached DF = DAG runs twice.
Подвох Action перезапускает всю родословную каждый раз, если не сделать cache(). count(), а потом write() на некешированном DF = DAG выполняется дважды.
Interviewer will probe "Where does the shuffle physically live?" → executor local disk (shuffle data + index files), served by the external shuffle service so it survives executor loss. "What's a stage boundary?" → a wide dependency.
О чём спросит интервьюер «Где физически лежит shuffle?» → локальный диск executor (данные shuffle + индекс-файлы), обслуживается external shuffle service, чтобы пережить потерю executor. «Что такое граница стадии?» → широкая зависимость.

4 · Narrow vs Wide Dependencies & Shuffles4 · Узкие vs широкие зависимости и shuffle

NARROW (pipelined, no network) WIDE (shuffle, stage boundary) P0 ─► C0 map/filter/union P0 ─┐ P1 ─► C1 coalesce(down) P1 ─┼─► C0 groupByKey / join / P2 ─► C2 P2 ─┘ distinct / repartition 1 parent → at most 1 child many parents → 1 child
УЗКАЯ (конвейерная, без сети) ШИРОКАЯ (shuffle, граница стадии) P0 ─► C0 map/filter/union P0 ─┐ P1 ─► C1 coalesce(вниз) P1 ─┼─► C0 groupByKey / join / P2 ─► C2 P2 ─┘ distinct / repartition 1 родитель → макс. 1 ребёнок много родителей → 1 ребёнок
NarrowWide (shuffle)
NetworkNoneDisk write + network fetch
Examplesmap, filter, union, mapPartitionsgroupByKey, reduceByKey, join*, distinct, repartition
CostCheap, pipelinedDominant cost + skew surface
RecoveryRecompute 1 partitionMay recompute upstream
УзкаяШирокая (shuffle)
СетьНетЗапись на диск + fetch по сети
Примерыmap, filter, union, mapPartitionsgroupByKey, reduceByKey, join*, distinct, repartition
ЦенаДёшево, конвейерноДоминирующая стоимость + площадка для перекоса
ВосстановлениеПеревычислить 1 партициюМожет перевычислить upstream
Probe "reduceByKey vs groupByKey?" → reduceByKey does map-side combine before the shuffle → far less data moved. groupByKey shuffles everything → OOM on hot keys.
Вопрос «reduceByKey vs groupByKey?» → reduceByKey делает комбинацию на стороне map перед shuffle → гораздо меньше данных. groupByKey шафлит всё → OOM на горячих ключах.
Minimize shuffles: broadcast small joins · filter/aggregate before shuffling · reduceByKey over groupByKey · bucket on join keys · prune columns · right-size partitions.
Минимизировать shuffle: broadcast небольших джойнов · фильтровать/агрегировать до shuffle · reduceByKey вместо groupByKey · bucket по ключам джойна · обрезать колонки · правильно подбирать партиции.
Anchor — Front Tier "On my Scala/Spark training pipeline on Hadoop, shuffles on a categorical key were the bottleneck; I cut them with map-side aggregation and salting."
Якорь — Front Tier «В моём Scala/Spark пайплайне обучения на Hadoop shuffle по категориальному ключу были узким местом; я обрезал их через агрегацию на стороне map и salting.»

5 · Join Strategies5 · Стратегии джойнов

StrategyWhenMechanism
Broadcast Hash (BHJ)One side < ~10MB (autoBroadcastJoinThreshold)Ship small side to all executors, build hash map. No big-side shuffle. Fastest.
Sort-Merge (SMJ)Default large–large equi-joinShuffle both by key, sort, merge. Robust, pays full shuffle+sort.
Shuffle Hash (SHJ)One side fits per-partition; sort too costlyHash one side per partition, no sort.
Broadcast Nested Loop / CartesianNon-equi / crossFallback; expensive.
СтратегияКогдаМеханизм
Broadcast Hash (BHJ)Одна сторона < ~10MB (autoBroadcastJoinThreshold)Отправить маленькую сторону всем executors, построить хеш-карту. Без shuffle большой стороны. Быстрейший.
Sort-Merge (SMJ)По умолчанию для больших equi-joinShuffle обеих по ключу, сортировка, слияние. Надёжный, платит полный shuffle+сортировка.
Shuffle Hash (SHJ)Одна сторона влезает на партицию; сортировка слишком дорогаяХеш одной стороны на партицию, без сортировки.
Broadcast Nested Loop / CartesianNon-equi / crossРезервный; дорогой.
# force broadcast (PySpark)
from pyspark.sql.functions import broadcast
big.join(broadcast(small), "key")
-- SQL hint
SELECT /*+ BROADCAST(dim) */ ... FROM fact JOIN dim ...
# принудительный broadcast (PySpark)
from pyspark.sql.functions import broadcast
big.join(broadcast(small), "key")
-- SQL hint
SELECT /*+ BROADCAST(dim) */ ... FROM fact JOIN dim ...
Probe "When does broadcast hurt?" → side bigger than judged → driver OOM gathering it or executor OOM holding the map; many concurrent broadcasts. Disable with threshold -1.
Вопрос «Когда broadcast вредит?» → сторона больше, чем посчитано → OOM драйвера при сборе или OOM executor при хранении карты; много одновременных broadcasts. Отключить порогом -1.
Current-knowledge flex AQE can convert SMJ → BHJ at runtime once it sees real post-filter sizes.
Демонстрация текущих знаний AQE может конвертировать SMJ → BHJ во время выполнения, увидев реальные размеры после фильтров.

6 · Catalyst Optimizer & Tungsten6 · Оптимизатор Catalyst и Tungsten

Catalyst = planning

Catalyst = планирование

Extensible rule-based + cost-based optimizer over plan trees.

Расширяемый rule-based + cost-based оптимизатор над деревьями планов.

  • Phases: analysis → logical opt → physical planning → codegen
  • Rules: predicate pushdown, column pruning, constant folding, filter reorder, partition pruning, join reordering (CBO)
  • CBO needs stats: ANALYZE TABLE … COMPUTE STATISTICS
  • Фазы: анализ → логическая опт. → физическое планирование → codegen
  • Правила: predicate pushdown, обрезка колонок, сворачивание констант, переупорядочивание фильтров, partition pruning, переупорядочивание джойнов (CBO)
  • CBO нужна статистика: ANALYZE TABLE … COMPUTE STATISTICS

Tungsten = execution

Tungsten = выполнение

  • Off-heap binary memory (no JVM object overhead/GC)
  • Whole-stage codegen: fuse operator chain into one generated Java function
  • Cache-aware + vectorized/columnar Parquet reads
  • Off-heap бинарная память (без оверхеда JVM-объектов/GC)
  • Whole-stage codegen: сливает цепочку операторов в одну сгенерированную Java-функцию
  • Cache-aware + векторизованное/колоночное чтение Parquet

One-liner: "Catalyst decides what plan; Tungsten makes its execution fast and memory-lean."

В одну строку: «Catalyst решает, какой план; Tungsten делает его выполнение быстрым и экономным по памяти.»

Read the plan df.explain(True) / explain("formatted") → look for Exchange (=shuffle), BroadcastHashJoin, *(n) codegen stages, pushed filters.
Читать план df.explain(True) / explain("formatted") → искать Exchange (=shuffle), BroadcastHashJoin, *(n) стадии codegen, pushed-фильтры.

7 · Caching & Persistence7 · Кеширование и персистентность

  • cache() = persist(MEMORY_AND_DISK) for DataFrames (RDD cache = MEMORY_ONLY).
  • Levels: MEMORY_ONLY · MEMORY_AND_DISK · *_SER · DISK_ONLY · *_2 (replicated) · off-heap.
  • Cache when a DF is reused across multiple actions / iterations (ML loops). Don't cache a once-through pipeline.
  • cache() = persist(MEMORY_AND_DISK) для DataFrames (кеш RDD = MEMORY_ONLY).
  • Уровни: MEMORY_ONLY · MEMORY_AND_DISK · *_SER · DISK_ONLY · *_2 (реплицированный) · off-heap.
  • Кешировать, когда DF переиспользуется через несколько actions / итераций (ML-циклы). Не кешировать одноразовый пайплайн.
Gotchas Cache is lazy (materializes on next action) · over-caching evicts & forces recompute/spill · forgetting unpersist() leaks memory.
Подвохи Кеш — ленивый (материализуется на следующем action) · избыток кеша вытесняет и заставляет перевычислять/spill · забытый unpersist() течёт память.
Anchor ML training (Front Tier) = the canonical cache case: the feature matrix is iterated many times.
Якорь ML-обучение (Front Tier) = канонический случай кеша: матрица признаков итерируется много раз.

8 · Partitioning & Pruning8 · Партиционирование и обрезка

ConceptMeaning
Input partitionsFrom source (file/block size, maxPartitionBytes)
Shuffle partitionsspark.sql.shuffle.partitions (default 200) — AQE coalesces now
repartition(n[,col])Full shuffle; increase/even out; can hash by col
coalesce(n)No shuffle; only merges (reduce); can cause skew
write.partitionBy(col)On-disk directory layout for later pruning (different from repartition)
КонцепцияСмысл
Входные партицииИз источника (размер файла/блока, maxPartitionBytes)
Shuffle-партицииspark.sql.shuffle.partitions (по умолчанию 200) — AQE сливает теперь
repartition(n[,col])Полный shuffle; увеличить/выровнять; можно хешировать по col
coalesce(n)Без shuffle; только слияние (уменьшение); может вызвать перекос
write.partitionBy(col)На-диске структура папок для последующего pruning (отличается от repartition)
Probe "Reduce 2000→50 before write?" → coalesce(50) (no shuffle) unless skewed, then repartition(50).
Вопрос «Уменьшить 2000→50 перед записью?» → coalesce(50) (без shuffle), если только не перекос, тогда repartition(50).
Predicate pushdown: filter into the file format (Parquet row-group min/max skip).
Partition pruning: physically partitioned table (by ds/date) → read only matching dirs.
Dynamic Partition Pruning: prune the big fact table at runtime using keys from the small broadcast dim side — star-schema win.
Predicate pushdown: фильтр в формат файла (пропуск row-group Parquet по min/max).
Partition pruning: физически партиционированная таблица (по ds/дате) → читать только подходящие папки.
Dynamic Partition Pruning: обрезать большую fact-таблицу в runtime, используя ключи из маленькой broadcast dim-стороны — выигрыш для star-схемы.
Anchor — Meta Date-partitioned (ds) tables, reading only relevant partitions = daily bread; same idea as Hive/Iceberg pruning.
Якорь — Meta Таблицы, партиционированные по дате (ds), чтение только релевантных партиций = ежедневный хлеб; та же идея, что Hive/Iceberg pruning.

9 · Skew Handling & AQE9 · Работа с перекосом и AQE

AQE — Adaptive Query Execution (on by default 3.2+)

AQE — Adaptive Query Execution (по умолчанию с 3.2+)

Re-optimizes the physical plan at runtime from real shuffle stats:

Переоптимизирует физический план в runtime из реальной статистики shuffle:

  • Coalesce shuffle partitions — collapse tiny post-shuffle partitions (kills the "200 default, mostly empty" problem)
  • Switch join strategy — SMJ → BHJ when a side is small after filters
  • Optimize skew joins — split skewed partitions into sub-partitions
  • Слияние shuffle-партиций — сворачивает крошечные партиции после shuffle (убивает проблему «200 по умолчанию, в основном пустых»)
  • Переключение стратегии джойна — SMJ → BHJ, когда сторона маленькая после фильтров
  • Оптимизация перекоса в джойнах — разбивает перекошенные партиции на подпартиции
Probe "What can't AQE fix?" → bad layout / no pruning, fundamentally too much data shuffled, opaque UDF plans, skew inside a single partition's compute.
Вопрос «Что AQE не может исправить?» → плохая структура / нет pruning, фундаментально слишком много данных в shuffle, непрозрачные UDF-планы, перекос внутри вычисления одной партиции.

Skew toolkit (concrete)

Инструментарий для перекоса (конкретно)

  1. AQE skew join — first line of defense.
  2. Salt the hot key: suffix 0..N on both sides (replicate small side across salts), join on salted key, aggregate back.
  3. Broadcast the small side → remove the shuffle entirely.
  4. Isolate hot keys: process top-N separately (often broadcast), union the rest.
  5. Pre-aggregate (map-side combine) before the shuffle.
  6. Handle nulls: null join keys all hash to one partition — filter/coalesce them.
  1. AQE skew join — первая линия защиты.
  2. Соление (salt) горячего ключа: добавить суффикс 0..N с обеих сторон (реплицировать маленькую сторону по солям), джойнить по солёному ключу, агрегировать обратно.
  3. Broadcast маленькой стороны → убрать shuffle полностью.
  4. Изолировать горячие ключи: обрабатывать топ-N отдельно (часто broadcast), объединить остальное.
  5. Предагрегация (комбинация на стороне map) перед shuffle.
  6. Обработать null: null join-ключи все хешируются в одну партицию — фильтровать/сливать их.
Probe "How do you detect skew?" → Spark UI task duration / shuffle-read max-vs-median; partition record counts. 199/200 done = textbook skew.
Вопрос «Как обнаружить перекос?» → Spark UI продолжительность задач / shuffle-read макс. vs медиана; количество записей на партицию. 199/200 готово = учебник перекоса.

10 · Executor Memory Model10 · Модель памяти Executor

Executor JVM heap ┌────────────────────────────────────────────┐ │ Reserved ~300MB │ ├────────────────────────────────────────────┤ │ Unified (spark.memory.fraction ≈ 0.6) │ │ ┌─────────────┬──────────────────────┐ │ │ │ Execution │ Storage (cache) │ │ ← borrow each other │ │ shuffle/sort│ cached blocks │ │ dynamically │ │ /agg/join │ │ │ │ └─────────────┴──────────────────────┘ │ ├────────────────────────────────────────────┤ │ User memory ≈ 0.4 (UDF state, your structs) │ └────────────────────────────────────────────┘ + Off-heap (Tungsten) + memoryOverhead (PySpark workers, netty)
Куча JVM Executor ┌────────────────────────────────────────────┐ │ Зарезервированная ~300MB │ ├────────────────────────────────────────────┤ │ Unified (spark.memory.fraction ≈ 0.6) │ │ ┌─────────────┬──────────────────────┐ │ │ │ Execution │ Storage (кеш) │ │ ← берут друг у друга │ │ shuffle/sort│ кешированные блоки │ │ динамически │ │ /agg/join │ │ │ │ └─────────────┴──────────────────────┘ │ ├────────────────────────────────────────────┤ │ Пользовательская память ≈ 0.4 (состояние UDF, твои структуры) │ └────────────────────────────────────────────┘ + Off-heap (Tungsten) + memoryOverhead (PySpark workers, netty)
  • Execution can evict Storage beyond a floor; the reverse is bounded (asymmetric).
  • Spill: execution memory exhausted → sort/agg buffers go to disk (slower, avoids OOM).
  • OOM causes: huge broadcast · skewed partition too big · collect() to driver · too few partitions.
  • Execution может вытеснить Storage сверх минимума; обратное ограничено (асимметрично).
  • Spill: execution-память исчерпана → буферы сортировки/агрегации уходят на диск (медленнее, избегает OOM).
  • Причины OOM: огромный broadcast · перекошенная партиция слишком большая · collect() на драйвер · слишком мало партиций.
memoryOverhead gotcha PySpark workers + off-heap buffers live outside executor.memory → tune spark.executor.memoryOverhead when PySpark OOMs mysteriously.
Подвох memoryOverhead PySpark-воркеры + off-heap буферы живут вне executor.memory → настраивай spark.executor.memoryOverhead, когда PySpark загадочно OOMится.

11 · Spark on YARN / Kubernetes & Fault Tolerance11 · Spark на YARN / Kubernetes и отказоустойчивость

Cluster managers & modes

Менеджеры кластера и режимы

  • YARN — Hadoop-native, shared queues/capacity scheduler (Front Tier world)
  • Kubernetes — executors as pods, cloud-native isolation, pod templates
  • Standalone / Mesos (legacy)
  • client mode = driver on submit host (interactive); cluster mode = driver in cluster (prod)
  • Dynamic allocation scales executors by pending tasks; needs external shuffle service (or k8s shuffle tracking)
  • YARN — нативный для Hadoop, общие очереди/capacity scheduler (мир Front Tier)
  • Kubernetes — executors как поды, cloud-native изоляция, шаблоны подов
  • Standalone / Mesos (легаси)
  • client режим = драйвер на хосте подачи (интерактивный); cluster режим = драйвер в кластере (прод)
  • Dynamic allocation масштабирует executors по pending-задачам; нужен external shuffle service (или k8s shuffle tracking)

Fault tolerance

Отказоустойчивость

  • Lineage (DAG): recompute lost partitions vs replicate data
  • Shuffle files on disk + external shuffle service survive executor loss
  • Checkpointing: truncate long lineage to reliable storage (HDFS); required for streaming state
  • Task retries + speculative execution for stragglers
  • Driver = SPOF → cluster mode + supervision / streaming checkpoints
  • Родословная (DAG): перевычислить потерянные партиции vs реплицировать данные
  • Shuffle-файлы на диске + external shuffle service переживают потерю executor
  • Checkpointing: обрезать длинную родословную в надёжное хранилище (HDFS); требуется для streaming state
  • Ретраи задач + speculative execution для отстающих
  • Драйвер = SPOF → cluster mode + supervision / streaming checkpoints

12 · Structured Streaming Basics12 · Основы Structured Streaming

  • Model: stream = unbounded table; each micro-batch appends rows; same DataFrame API, incremental execution. (Continuous mode = low-latency, experimental.)
  • Triggers: default (ASAP) · fixed interval · availableNow/once (batch backfill) · continuous
  • Output modes: append (new rows; needs watermark for aggs) · update · complete (whole result, aggs only)
  • Watermark: withWatermark("ts","10 minutes") bounds state for late data; later events dropped
  • State + checkpoint: stateful ops keep a state store, checkpointed → exactly-once with replayable source (Kafka offsets) + idempotent sink
  • Модель: поток = неограниченная таблица; каждый микро-батч добавляет строки; тот же API DataFrame, инкрементальное выполнение. (Continuous mode = низкая задержка, экспериментально.)
  • Триггеры: по умолчанию (ASAP) · фиксированный интервал · availableNow/once (батч-backfill) · continuous
  • Output modes: append (новые строки; нужен watermark для агрегатов) · update · complete (весь результат, только агрегаты)
  • Watermark: withWatermark("ts","10 minutes") ограничивает state для поздних данных; более поздние события отбрасываются
  • State + checkpoint: stateful-операции держат state store, чекпоинтятся → exactly-once с перечитываемым источником (Kafka offsets) + идемпотентным приёмником
Anchor — IU Group Kafka ingestion for Syntea (80k+ students) is your streaming story.
Якорь — IU Group Kafka-ingestion для Syntea (80k+ студентов) — твоя история про стриминг.
Gap-bridge Flink: true record-at-a-time, richer event-time/state, savepoints. Concepts (event time, watermark, window, state, exactly-once) transfer directly; API/runtime differ. Say this confidently.
Мост через пробел Flink: настоящий запись-за-раз, более богатое event-time/state, savepoints. Концепции (event time, watermark, window, state, exactly-once) переносятся напрямую; API/runtime отличаются. Говори это уверенно.

13 · Performance Tuning & Debugging Slow Jobs13 · Тюнинг производительности и отладка медленных джобов

The signature DE question: "this stage hangs at 199/200"

Фирменный DE-вопрос: «эта стадия зависла на 199/200»

  1. Spark UI first — stage timeline + task metrics. 199/200 = classic skew.
  2. Confirm skew: max vs median shuffle-read / duration (50x = skew).
  3. Rule out: skewed key (null/hot) · too few partitions · spill · GC · straggler node · small-files explosion.
  4. Fix by cause: AQE skew / salt / broadcast · filter nulls · raise partitions · more memory · speculative execution.
  5. Verify in the UI: compare shuffle-read distribution & stage time.
  1. Spark UI сначала — timeline стадии + метрики задач. 199/200 = классический перекос.
  2. Подтвердить перекос: макс. vs медиана shuffle-read / продолжительность (50x = перекос).
  3. Исключить: перекошенный ключ (null/горячий) · слишком мало партиций · spill · GC · отстающий узел · взрыв мелких файлов.
  4. Исправить по причине: AQE skew / salt / broadcast · фильтровать null · увеличить партиции · больше памяти · speculative execution.
  5. Проверить в UI: сравнить распределение shuffle-read и время стадии.

General tuning checklist

Общий чек-лист тюнинга

  • Read less: column/predicate/partition pruning · columnar Parquet/ORC · compact small files
  • Shuffle less: broadcast small joins · pre-aggregate · bucket on join key · avoid needless repartition
  • Right-size parallelism: partitions ≈ 2–4× cores · let AQE coalesce
  • Fix skew · cache only reused data
  • Executor sizing: ~5 cores/executor heuristic (fat executors → GC/HDFS throughput pain) · memoryOverhead · dynamic allocation
  • Avoid Python UDFs (serialization boundary) → built-in SQL fns or Arrow/pandas vectorized UDFs
  • Measure before/after — never tune blind
  • Читать меньше: обрезка колонок/предикатов/партиций · колоночный Parquet/ORC · компактировать мелкие файлы
  • Shuffle меньше: broadcast маленьких джойнов · предагрегация · bucket по ключу джойна · избегать ненужных repartition
  • Подобрать параллелизм: партиции ≈ 2–4× ядра · пусть AQE сливает
  • Исправить перекос · кешировать только переиспользуемые данные
  • Размер Executor: эвристика ~5 ядер/executor (жирные executors → боль GC/пропускная способность HDFS) · memoryOverhead · dynamic allocation
  • Избегать Python UDFs (граница сериализации) → встроенные SQL-функции или Arrow/pandas векторизованные UDFs
  • Измерять до/после — никогда не тюнить вслепую

15 · Self-Quiz — flip to reveal15 · Самопроверка — переверни, чтобы открыть

Tap a card to flip. Answer out loud first.

Нажми на карточку, чтобы перевернуть. Сначала ответь вслух.

What defines a stage boundary?Что определяет границу стадии?
tap to revealнажми, чтобы открыть
A wide dependency (shuffle). Each stage is a pipeline of narrow transforms; the shuffle splits stages.Широкая зависимость (shuffle). Каждая стадия — конвейер узких трансформаций; shuffle разделяет стадии.
reduceByKey vs groupByKey?reduceByKey vs groupByKey?
tap to revealнажми, чтобы открыть
reduceByKey = map-side combine before shuffle → less data. groupByKey shuffles all values → OOM risk on hot keys.reduceByKey = комбинация на стороне map перед shuffle → меньше данных. groupByKey шафлит все значения → риск OOM на горячих ключах.
Catalyst vs Tungsten?Catalyst vs Tungsten?
tap to revealнажми, чтобы открыть
Catalyst = planning/optimizer (what plan). Tungsten = execution engine (off-heap memory, whole-stage codegen).Catalyst = планирование/оптимизатор (какой план). Tungsten = движок выполнения (off-heap память, whole-stage codegen).
3 things AQE does at runtime?3 вещи, которые AQE делает в runtime?
tap to revealнажми, чтобы открыть
1) Coalesce shuffle partitions · 2) SMJ→BHJ join switch · 3) Split skewed partitions. On by default 3.2+.1) Слияние shuffle-партиций · 2) Переключение джойна SMJ→BHJ · 3) Разделение перекошенных партиций. По умолчанию с 3.2+.
repartition vs coalesce?repartition vs coalesce?
tap to revealнажми, чтобы открыть
repartition = full shuffle, up or down, can hash by col. coalesce = no shuffle, only reduce, can skew.repartition = полный shuffle, вверх или вниз, можно хешировать по колонке. coalesce = без shuffle, только уменьшение, может перекосить.
Stage hangs at 199/200 — what is it?Стадия зависла на 199/200 — что это?
tap to revealнажми, чтобы открыть
Data skew. One task got a hot/null key. Fix: AQE skew join, salting, broadcast, isolate hot keys, filter nulls.Перекос данных. Одна задача получила горячий/null ключ. Лечение: AQE skew join, соление, broadcast, изолировать горячие ключи, фильтровать null.
Default broadcast threshold?Порог broadcast по умолчанию?
tap to revealнажми, чтобы открыть
10MB (spark.sql.autoBroadcastJoinThreshold). Disable with -1. Force with broadcast(df).10MB (spark.sql.autoBroadcastJoinThreshold). Отключить через -1. Принудительно через broadcast(df).
Exactly-once in Structured Streaming needs?Exactly-once в Structured Streaming требует?
tap to revealнажми, чтобы открыть
Replayable source (Kafka offsets) + checkpointing + idempotent sink.Перечитываемый источник (Kafka offsets) + чекпоинтинг + идемпотентный приёмник.
What is a watermark for?Для чего watermark?
tap to revealнажми, чтобы открыть
Bounds how long streaming state is kept for late data so it doesn't grow forever; events later than it are dropped.Ограничивает, как долго streaming state хранится для поздних данных, чтобы не рос вечно; события позже него отбрасываются.
Why prefer Dataset over PySpark DataFrame?Почему предпочесть Dataset вместо PySpark DataFrame?
tap to revealнажми, чтобы открыть
JVM-native (no Python ser boundary) + compile-time type safety → testable, maintainable, refactorable.JVM-нативный (без границы сериализации Python) + типобезопасность на этапе компиляции → тестируемый, поддерживаемый, рефакторимый.