YARN: Resource Manager, Node Manager и режимы запуска Spark

YARN - стандартный Cluster Manager в Hadoop-экосистеме. Разбираем архитектуру, два режима spark-submit (client vs cluster) и типичные проблемы в enterprise-кластерах.

platform

Архитектура YARN

  • Resource Manager (RM) - глобальный планировщик. Знает ресурсы всего кластера, принимает заявки от ApplicationMaster
  • Node Manager (NM) - агент на каждом Worker-узле. Запускает Containers, мониторит использование ресурсов, отчитывается RM
  • ApplicationMaster (AM) - процесс-координатор конкретного приложения Spark (Driver в cluster mode или отдельный процесс в client mode). Запрашивает Containers у RM для Executor'ов
  • Container - набор ресурсов (CPU + RAM) на Node Manager'е. Каждый Executor = один Container

Два режима запуска: client vs cluster

Client mode (по умолчанию)

spark-submit \
  --master yarn \
  --deploy-mode client \   # Driver работает на машине, откуда запустили submit
  --num-executors 10 \
  --executor-memory 4g \
  --executor-cores 4 \
  job.py
[Laptop / Edge Node]          [YARN Cluster]
  Driver Process    ←→  AM   ←→  RM
  (ваш процесс)              ↓
                           Executors (Containers)
  • Driver работает на машине откуда запущен submit (ноутбук, edge node)
  • Если соединение прервётся - Driver умирает → Job падает
  • Spark UI доступен на локальной машине (localhost:4040)
  • Используется для разработки и отладки

Cluster mode

spark-submit \
  --master yarn \
  --deploy-mode cluster \  # Driver запускается как Container внутри YARN
  --num-executors 10 \
  --executor-memory 4g \
  --executor-cores 4 \
  s3://bucket/job.py
[Laptop]             [YARN Cluster]
  submit ──────→    AM = Driver (Container)  ←→  RM
  (только запуск)                                ↓
                                          Executors (Containers)
  • Driver запускается как Container внутри кластера
  • Процесс submit завершается сразу после отправки приложения
  • Job продолжается независимо от клиентской машины
  • Рекомендуется для production (нет зависимости от подключения)

Ключевые параметры YARN для Spark

spark-submit \
  --master yarn \
  --deploy-mode cluster \
  --driver-memory 4g \               # память Driver Container
  --driver-cores 2 \                 # CPU Driver Container
  --num-executors 20 \               # число Executor Containers
  --executor-memory 8g \             # RAM Executor Container
  --executor-cores 4 \               # CPU Executor Container
  --conf spark.yarn.executor.memoryOverhead=1g \   # overhead для Python workers
  --conf spark.yarn.driver.memoryOverhead=512m \
  --queue default \                  # YARN queue (Fair/Capacity Scheduler)
  job.py

YARN Queues (Fair/Capacity Scheduler)

# Посмотреть доступные queues
yarn queue -status default
yarn queue -status data-engineering

# Задать queue для Spark job
spark-submit --queue data-engineering job.py

# Приоритет в Fair Scheduler - через конфиг
spark.conf.set("spark.yarn.scheduler.reporterThread.maxFailures", "5")

Часто встречаемые ошибки YARN

Ошибка Причина Решение
Application is added to the scheduler Нет ресурсов в очереди Ждать или увеличить quota queue
Container killed by YARN for exceeding memory limits executor.memory + overhead > Container limit Увеличить spark.yarn.executor.memoryOverhead
Application is Killed YARN timeout или quota exceeded Проверить yarn application -status <app_id>
AM container is killed Driver OOM Увеличить --driver-memory
# Диагностика через YARN CLI
yarn application -list                           # текущие приложения
yarn application -status application_XXX_0001   # статус конкретного
yarn logs -applicationId application_XXX_0001   # логи Driver + Executors
yarn node -list                                  # состояние узлов
yarn node -status <nodeId>                       # детали узла

Dynamic Resource Allocation на YARN

spark.conf.set("spark.dynamicAllocation.enabled",            "true")
spark.conf.set("spark.shuffle.service.enabled",              "true")   # обязательно для YARN
spark.conf.set("spark.dynamicAllocation.minExecutors",        "2")
spark.conf.set("spark.dynamicAllocation.maxExecutors",        "50")
spark.conf.set("spark.dynamicAllocation.initialExecutors",    "5")
spark.conf.set("spark.dynamicAllocation.executorIdleTimeout", "60s")

spark.shuffle.service.enabled = true обязателен: ExternalShuffleService работает как daemon на каждом Node Manager'е и хранит shuffle файлы независимо от Executor'ов - позволяет их убивать при простое.