Ответ
После выполнения 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) внутри драйвера:
- Разбивает вычисления, построенные на RDD или DataFrame, на Directed Acyclic Graph (DAG) операций.
- Разделяет DAG на Stages (этапы) на основе операций shuffle (например,
groupBy,join). - Каждый 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.