DAG-планировщик Spark: от логического плана к стадиям

Когда вы отправляете задание в Spark, DAG-планировщик переводит логический план — цепочку преобразований RDD — в физический план выполнения:

  1. Построить DAG из последовательности преобразований
  2. Найти границы стадий в точках shuffle, то есть на широких зависимостях вроде groupByKey и reduceByKey
  3. Топологически отсортировать стадии, чтобы определить порядок выполнения
  4. Отправить задачи по партициям каждой стадии в планировщик задач

Стадии внутри одного конвейера, связанные узкими зависимостями, сливают свои преобразования и обходятся без промежуточной материализации.

DAG из десяти RDD и action, с двумя join: планировщик идёт назад от action, режет граф на стадии по четырём широким (shuffle) зависимостям, сортирует стадии топологически и запускает их волнами, а нажатие на стрелку делает её узкой или широкой.

Пройдите план по шагам или нажмите на стрелку: узкая зависимость станет shuffle (а сторона join — заранее распределённой по ключу), и стадии перестроятся, а волны изменятся. Нажатие на RDD покажет его стадию и ту зависимость, из-за которой граница проходит именно там.

Почему граница стадии — это именно shuffle

Узкая зависимость означает, что каждая партиция-потомок читает ровно одну партицию-предка. Такую цепочку можно выполнить на одном исполнителе, не выходя в сеть, и Spark сливает её в один конвейер: map, filter и mapValues идут подряд над одной записью, не записывая промежуточных результатов.

Широкая зависимость означает, что партиция-потомок читает многие партиции-предки. Пока все предки не досчитаны и не записаны, ни одна партиция-потомок не может начаться — это барьер. Именно поэтому граница стадии совпадает с shuffle: стадия — это максимальный кусок работы, который не требует барьера внутри себя.

Отсюда следует практическое правило: число стадий в задании равно числу shuffle плюс один. Если в Spark UI стадий больше, чем вы ожидали, значит где-то в плане есть shuffle, о котором вы не думали — чаще всего это distinct, join по неотпартиционированному ключу или repartition.