Что происходит после отправки задания в Spark (spark-submit)?

«Что происходит после отправки задания в Spark (spark-submit)?» — вопрос из категории Apache Spark, который задают на 33% собеседований Data Инженер. Ниже — развёрнутый ответ с разбором ключевых моментов.

Ответ

После выполнения spark-submit запускается цепочка событий для распределенного выполнения кода. Вот что происходит под капотом на примере кластера YARN:

1. Запуск драйвера (Driver Program):

  • Команда spark-submit инициирует запуск Spark Driver — главного процесса вашего приложения.
  • В режиме --deploy-mode client драйвер работает на той же машине, откуда была запущена команда. В режиме cluster — на одной из нод кластера (обычно на ApplicationMaster в YARN).
  • Драйвер создает SparkContext, который является точкой входа во все функции Spark.

2. Запрос ресурсов у кластера:

  • SparkContext подключается к менеджеру ресурсов кластера (YARN ResourceManager, Mesos Master или Spark Standalone Master).
  • Он запрашивает ресурсы (ядра CPU и память) для Executor'ов — рабочих процессов, которые будут выполнять задачи.

3. Распределение и выполнение кода:

  • Менеджер ресурсов выделяет контейнеры на Worker-нодах и запускает в них Executor'ы.
  • Driver отправляет код приложения (JAR или Python-файлы) и задачи на Executor'ы.
  • Планировщик (DAGScheduler) внутри драйвера:
    1. Разбивает вычисления, построенные на RDD или DataFrame, на Directed Acyclic Graph (DAG) операций.
    2. Разделяет DAG на Stages (этапы) на основе операций shuffle (например, groupBy, join).
    3. Каждый Stage состоит из множества Tasks — элементарных единиц работы, которые выполняют одну и ту же функцию над разными частями данных (партициями).
  • TaskScheduler отправляет Tasks на свободные Executor'ы для параллельного выполнения.

4. Мониторинг и завершение:

  • Driver отслеживает выполнение Tasks, собирает результаты (например, при действиях collect(), count()) или записывает данные во внешние системы.
  • После завершения всех задач или при вызове spark.stop() Spark освобождает ресурсы, и приложение завершается.

Пример и наблюдение:

spark-submit 
  --master yarn 
  --deploy-mode cluster 
  --executor-memory 4g 
  my_etl_job.py

Ход выполнения можно отслеживать через веб-UI Spark на порту 4040 драйвера или через интерфейс YARN ResourceManager.