DAG-планировщик Spark: от логического плана к стадиям
Когда вы отправляете задание в Spark, DAG-планировщик переводит логический план — цепочку преобразований RDD — в физический план выполнения:
- Построить DAG из последовательности преобразований
- Найти границы стадий в точках shuffle, то есть на широких зависимостях вроде
groupByKeyиreduceByKey - Топологически отсортировать стадии, чтобы определить порядок выполнения
- Отправить задачи по партициям каждой стадии в планировщик задач
Стадии внутри одного конвейера, связанные узкими зависимостями, сливают свои преобразования и обходятся без промежуточной материализации.
DAG из десяти RDD и action, с двумя join: планировщик идёт назад от action, режет граф на стадии по четырём широким (shuffle) зависимостям, сортирует стадии топологически и запускает их волнами, а нажатие на стрелку делает её узкой или широкой.
Пройдите план по шагам или нажмите на стрелку: узкая зависимость станет shuffle (а сторона join — заранее распределённой по ключу), и стадии перестроятся, а волны изменятся. Нажатие на RDD покажет его стадию и ту зависимость, из-за которой граница проходит именно там.
Почему граница стадии — это именно shuffle
Узкая зависимость означает, что каждая партиция-потомок читает ровно одну партицию-предка. Такую цепочку можно выполнить на одном исполнителе, не выходя в сеть, и Spark сливает её в один конвейер: map, filter и mapValues идут подряд над одной записью, не записывая промежуточных результатов.
Широкая зависимость означает, что партиция-потомок читает многие партиции-предки. Пока все предки не досчитаны и не записаны, ни одна партиция-потомок не может начаться — это барьер. Именно поэтому граница стадии совпадает с shuffle: стадия — это максимальный кусок работы, который не требует барьера внутри себя.
Отсюда следует практическое правило: число стадий в задании равно числу shuffle плюс один. Если в Spark UI стадий больше, чем вы ожидали, значит где-то в плане есть shuffle, о котором вы не думали — чаще всего это distinct, join по неотпартиционированному ключу или repartition.