Block Locality: как Spark scheduler использует data-locality при планировании

Фундаментальный разбор Data Locality в HDFS-кластерах: физика блочного хранения, иерархия уровней PROCESS_LOCAL → ANY, алгоритм Delay Scheduling, механизм Short-Circuit Reads, влияние Shuffle и кеширования, мониторинг через Spark UI и тюнинг под SLA.

core storage

1. Физическая концепция Data Locality в on-premise архитектурах

Data Locality - это принцип, ставший одним из краеугольных камней Hadoop при его появлении в середине 2000-х. Понять его значение проще всего через сравнение двух противоположных архитектурных решений, каждое из которых отвечает на вопрос: «у нас есть данные и вычислительная мощность - как их свести вместе?»

Философия: код дешевле данных

В традиционных распределённых системах конца 1990-х - начала 2000-х данные хранились на выделенных серверах хранения, а вычислительные узлы получали их по сети. Логика была простой: серверы хранения дешевле вычислительных, можно их наполнить дисками и раздавать данные.

Проблема обнаружилась с ростом объёмов. Когда данных стало сотни терабайт, а потом петабайты, сеть превратилась в абсолютное узкое место. Прочитать 1 TB с дисков сервера - это несколько минут. Передать 1 TB по гигабитному каналу - это 2+ часа. По десяти-гигабитному - около 13 минут. Даже с отличной сетью передача данных через неё оказывается в разы медленнее их локального чтения.

Hadoop перевернул эту логику: вместо того чтобы приносить данные к программе, принести программу к данным. Код MapReduce (или Spark Task) занимает несколько сотен килобайт - JAR-файл. Данные - сотни мегабайт на блок. Очевидно, что перемещать маленький код к большим данным в сотни раз эффективнее, чем наоборот.

Именно поэтому в Hadoop-кластере каждый узел является одновременно и DataNode (хранит блоки HDFS), и NodeManager (запускает YARN-контейнеры с вычислениями). Compute и Storage совмещены физически на одном железе.

Схема наглядно показывает принципиальную разницу: в облаке данные всегда идут по сети между независимыми системами, тогда как в on-premise HDFS-кластере идеальный сценарий - чтение с локального диска без сетевого трафика вообще.

Анатомия блока HDFS: от файла до физического диска

Когда Spark читает файл hdfs://cluster/data/events/part-0000.parquet, этот файл не хранится как единый объект на одном диске. HDFS разбивает его на блоки фиксированного размера и распределяет реплики по нескольким DataNode.

Стандартный размер блока HDFS - 128 MB. Для аналитических нагрузок Data Lake часто используют 256 MB или 512 MB: меньше блоков = меньше обращений к NameNode = меньше overhead при планировании. Слишком большой блок (> 512 MB) увеличивает time-to-first-byte и снижает параллелизм при чтении.

Rack-aware репликация - HDFS по умолчанию размещает реплики не случайно:

  • Реплика 1: на том же узле, что записывает клиент (или на произвольном узле, если клиент - не DataNode)
  • Реплика 2: на другом узле в том же rack
  • Реплика 3: на узле в другом rack

Это обеспечивает одновременно и производительность (одна реплика всегда рядом), и отказоустойчивость (при потере целого rack данные остаются в другом).

Связка YARN + HDFS: как Executor'ы попадают на нужные узлы

Когда Spark-приложение запускается на YARN, Driver сначала обращается к NameNode за информацией о расположении блоков, затем передаёт эти предпочтения ResourceManager YARN.

Ключевой момент: ResourceManager не гарантирует размещение Executor'ов именно там, где лежат данные. Если на Server1 нет свободных ресурсов, YARN разместит Executor на другом сервере. Именно поэтому в Spark существует механизм Delay Scheduling - ожидание освобождения «правильных» ресурсов.


2. Спектр уровней локальности Spark (Locality Levels)

Spark определяет несколько уровней локальности для каждой задачи. Уровни упорядочены от лучшего к худшему. При планировании Spark пытается назначить задаче наиболее высокий достижимый уровень.

Типичная скорость чтения одного блока 128 MB:

  • PROCESS_LOCAL: менее 1 мс (данные уже в Java heap)
  • NODE_LOCAL: 0.25–0.65 сек (скорость локального диска)
  • RACK_LOCAL: 0.1–1 сек (зависит от загруженности ToR коммутатора)
  • ANY: 0.1–13+ сек (зависит от топологии и загруженности сети)

PROCESS_LOCAL: данные прямо в памяти JVM

Это наилучший из возможных сценариев. Он возникает, когда RDD или DataFrame были закешированы в памяти Executor через .cache() или .persist(StorageLevel.MEMORY_ONLY).

При повторном чтении закешированного датафрейма Spark знает, в каком именно Executor (JVM-процесс) и на каком разделе (partition) хранится каждый кусок данных. Если задача назначена на тот же Executor, который хранит нужный кешированный раздел - это PROCESS_LOCAL. Никакого I/O нет вообще: данные читаются из Java heap.

from pyspark.sql import SparkSession
from pyspark import StorageLevel

spark = SparkSession.builder.getOrCreate()

# Первое чтение: NODE_LOCAL (данные с HDFS)
df = spark.read.parquet("hdfs://cluster/data/events/")

# Кешируем в памяти Executor'ов.
# После этого повторные обращения к df дадут PROCESS_LOCAL
# для партиций, которые уже закешированы в конкретных Executor'ах.
df.cache()
df.count()  # триггерит материализацию кеша

# Теперь этот запрос выполнится с PROCESS_LOCAL для закешированных партиций
result = df.groupBy("event_type").count()
result.show()

PROCESS_LOCAL исчезает после unpersist(), перезапуска Executor'а или вытеснения из кеша при нехватке памяти.

NODE_LOCAL: локальный диск DataNode

Это наиболее распространённый и желанный уровень для аналитических Spark-джобов, читающих HDFS напрямую (без кеша). Executor запущен на том же физическом сервере, что и DataNode с нужным блоком.

Чтение происходит через local file system или через Unix Domain Socket (если настроен Short-Circuit Read - об этом в разделе 4). В обоих случаях данные не покидают пределы сервера и не используют сетевую карту.

Типичная скорость: на современных HDD HDFS-кластерах - 200–400 MB/s на блок. На SSD - 400–800 MB/s. Это принципиально выше, чем через сеть даже при 10 Gbps (практически достижимые ~900 MB/s для одного соединения при идеальных условиях).

RACK_LOCAL: та же стойка, один hop сети

Если на сервере с нужным блоком нет свободного Executor'а, следующий предпочтительный вариант - сервер в той же rack. Данные передаются через Top-of-Rack (ToR) коммутатор - один сетевой хоп.

Современные ToR коммутаторы работают на 10–25 Gbps. Практически достижимая скорость одного TCP-соединения - 800 MB/s–2.5 GB/s. Это заметно медленнее локального диска, но намного быстрее межстоечного трафика через core/aggregation коммутаторы.

Почему rack важен? В типичных дата-центрах стойки соединены ToR коммутаторами 10–25 Gbps uplink, а ToR между собой - через aggregation/core layer 40–100 Gbps, который делится между всеми стойками. В условиях активного Spark-чтения межстоечный bandwidth становится перегруженным ресурсом.

ANY: где угодно в кластере

Задача может быть запущена на любом свободном Executor'е в кластере. Блок данных читается по сети через полную топологию дата-центра. Это самый медленный вариант и единственный, при котором Spark-задача гарантированно создаёт межстоечный трафик.

ANY - это не обязательно проблема для небольших блоков или задач с малым объёмом данных. Проблема возникает, когда большинство задач в большом Stage работают на уровне ANY - тогда весь кластер начинает активно гонять данные по сети.

Как HDFS узнаёт топологию сети: rack awareness

Для правильной работы rack-aware репликации и Spark scheduler HDFS должен знать физическую топологию сети. Это настраивается через topology script - исполняемый файл, которому DataNode передаёт свой IP-адрес и получает в ответ rack ID.

#!/bin/bash
# /etc/hadoop/topology.sh
# Скрипт определения топологии стойки по IP-адресу.
# HDFS вызывает этот скрипт при регистрации каждого DataNode.

IP=$1  # IP адрес DataNode передаётся как аргумент

case $IP in
    # Rack 1: серверы 192.168.1.0/24
    192.168.1.*)  echo "/rack1" ;;
    # Rack 2: серверы 192.168.2.0/24
    192.168.2.*)  echo "/rack2" ;;
    # Rack 3: серверы 192.168.3.0/24
    192.168.3.*)  echo "/rack3" ;;
    # По умолчанию: все неизвестные в rack0
    *)            echo "/default-rack" ;;
esac
<!-- core-site.xml -->
<property>
  <name>net.topology.script.file.name</name>
  <value>/etc/hadoop/topology.sh</value>
</property>

<!--
  Количество IP-адресов, которое кешируется в топологии.
  Для кластеров > 100 нод увеличивайте это значение во избежание
  постоянных вызовов скрипта при heartbeat'ах DataNode.
-->
<property>
  <name>net.topology.script.number.of.cached.filenames</name>
  <value>100</value>
</property>

Без настроенного topology script все узлы попадают в /default-rack. HDFS всё равно будет реплицировать данные, но без rack-aware распределения: обе первые реплики могут оказаться в одной физической стойке, что снижает отказоустойчивость.


3. Механика работы Spark Scheduler: алгоритм Delay Scheduling

Знание о расположении блоков - это только полдела. Вопрос в том, что делать, когда «правильный» сервер занят. Это ключевая проблема, которую решает Delay Scheduling.

Проблема конкуренции за ресурсы

Представьте ситуацию: файл содержит 100 блоков, все они находятся на Server1 и Server2 (2 реплики каждого блока). Spark-джоба имеет 20 Executor'ов, распределённых по 10 серверам. Если блоков 100 и Executor'ов 20, но «правильных» серверов только 2 - многие задачи неизбежно будут ждать или читать данные по сети.

Реальная ситуация: у Server1 есть 16 ядер и 4 Executor'а, каждый с 4 ядрами. Каждый блок занимает одну таску. Значит, Server1 может параллельно обрабатывать 4 задачи. Если активных задач на Server1 больше 4 - планировщик сталкивается с выбором: ждать освобождения ядра на Server1 или немедленно запустить задачу на Server2, где данных нет, но ресурсы свободны.

Алгоритм Delay Scheduling: ожидание ради локальности

Delay Scheduling - алгоритм, предложенный Матеем Захарией (создатель Spark) в 2010 году. Идея проста: вместо немедленного запуска задачи на любом свободном ресурсе, планировщик ждёт некоторое время, надеясь, что освободится ресурс с нужным уровнем локальности.

Если за время ожидания нужный Executor не освобождается, планировщик снижает требования к локальности (Locality Level Downgrade):

  1. Ждём spark.locality.wait.process секунд для PROCESS_LOCAL → если нет, пробуем NODE_LOCAL
  2. Ждём spark.locality.wait.node секунд для NODE_LOCAL → если нет, пробуем RACK_LOCAL
  3. Ждём spark.locality.wait.rack секунд для RACK_LOCAL → если нет, запускаем на ANY

Это постепенное снижение, а не мгновенный переход к ANY. Логи планировщика при снижении уровня:

INFO TaskSetManager: Starting task 42.0 in stage 3.0 (TID 185, server4,
  executor 7, partition 42, RACK_LOCAL, 5102 bytes)

WARN TaskSetManager: Level for task 42 has been downgraded from RACK_LOCAL
  to ANY because it waited 3000ms without finding a RACK_LOCAL executor.

INFO TaskSetManager: Starting task 42.0 in stage 3.0 (TID 186, server7,
  executor 12, partition 42, ANY, 5102 bytes)

Параметры Delay Scheduling

from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("locality-tuning") \

    # spark.locality.wait - базовое время ожидания для всех уровней.
    # Применяется если не установлены специфичные параметры ниже.
    # Дефолт: 3 сек (3000мс). Значение "0" = немедленный переход к ANY.
    .config("spark.locality.wait", "3s") \

    # spark.locality.wait.process - время ожидания для PROCESS_LOCAL.
    # Имеет смысл только при работе с кешированными DataFrame.
    # Для джобов, читающих "холодные" данные с HDFS, установите в 0 —
    # PROCESS_LOCAL недостижим при первом чтении с диска.
    .config("spark.locality.wait.process", "5s") \

    # spark.locality.wait.node - время ожидания для NODE_LOCAL.
    # Наиболее важный параметр для HDFS-кластеров.
    # При загруженных нодах увеличьте до 5-10 секунд.
    # При коротких задачах (< 1 сек каждая) уменьшите —
    # ожидание 3 секунды больше самой задачи неэффективно.
    .config("spark.locality.wait.node", "10s") \

    # spark.locality.wait.rack - время ожидания для RACK_LOCAL.
    # Обычно достаточно дефолта. Снижайте если rack-local
    # сеть перегружена и ANY не намного хуже.
    .config("spark.locality.wait.rack", "3s") \

    .getOrCreate()

Когда увеличивать и когда уменьшать spark.locality.wait

Увеличивайте spark.locality.wait.node, если:

  • Задачи в Stage выполняются долго (> 10 сек каждая) - ожидание 3–10 сек незначительно
  • Кластер сильно загружен и NODE_LOCAL ресурсы освобождаются, просто нужно подождать
  • Сеть между стойками - узкое место (< 10 Gbps суммарно)

Уменьшайте (вплоть до 0), если:

  • Задачи короткие (< 1 сек каждая) - ожидание 3 сек больше самой задачи
  • Данные равномерно распределены по всему кластеру и все реплики одинаково доступны
  • Используете S3/MinIO вместо HDFS - там нет локальности по определению
  • Кластер почти пустой - NODE_LOCAL ресурсы всегда доступны без ожидания

Специальный случай - стриминговые джобы Structured Streaming с короткими микробатчами (< 5 сек). Установка spark.locality.wait = 0 критически важна: ожидание локальности может задержать микробатч больше его интервала.


4. Оптимизация на стыке HDFS и Spark: Short-Circuit Local Reads

NODE_LOCAL - хороший уровень локальности, но даже при «локальном» чтении в стандартной конфигурации HDFS есть неэффективность. Short-Circuit Reads - это оптимизация, устраняющая этот overhead.

Стандартный путь NODE_LOCAL чтения: скрытый оверхед

Когда Spark Executor (JVM процесс) хочет прочитать блок HDFS на том же сервере, стандартная процедура выглядит так:

Spark Executor (JVM)
    ↓
HDFS Client в JVM
    ↓ TCP-соединение (localhost:50010)
DataNode процесс (JVM)
    ↓ read() syscall
Файл на диске: /data/hdfs/current/blk_001234

Несмотря на то что оба процесса находятся на одном сервере, данные проходят через TCP-стек:

  1. DataNode читает блок с диска в свой Java heap
  2. DataNode сериализует данные в TCP-буфер ядра ОС
  3. Spark Executor читает из TCP-буфера в свой Java heap
  4. Данные оказываются скопированными дважды: disk → DataNode heap → Executor heap

Это «лишнее» копирование через TCP-стек loopback интерфейса - ненужный overhead CPU и памяти.

Short-Circuit Reads: прямой доступ к файлу блока

Short-Circuit (короткое замыкание) - механизм, позволяющий HDFS Client Spark'а обойти DataNode процесс и читать файл блока напрямую через файловую систему:

Данные при Short-Circuit Read копируются только один раз: disk → Executor heap. Нет TCP overhead, нет DataNode как посредника передачи данных.

Практический прирост производительности:

  • CPU utilization DataNode при чтении снижается на 20–40%
  • Пропускная способность чтения возрастает на 10–30%
  • Latency первого байта снижается на 15–25%

Конфигурация Short-Circuit Reads

<!-- hdfs-site.xml на всех DataNode и клиентских узлах (Spark) -->

<!--
  Включаем Short-Circuit Reads для HDFS клиентов.
  true: HDFS Client будет пытаться читать блоки напрямую если DataNode
  на том же хосте, иначе автоматически откатывается к TCP.
-->
<property>
  <name>dfs.client.read.shortcircuit</name>
  <value>true</value>
</property>

<!--
  Путь к Unix Domain Socket для коммуникации между HDFS Client и DataNode.
  DataNode создаёт этот сокет при старте.
  Путь должен существовать и быть доступен HDFS-пользователю.
-->
<property>
  <name>dfs.domain.socket.path</name>
  <value>/var/lib/hadoop-hdfs/dn_socket</value>
</property>

<!--
  Разрешить DataNode чтение файлов блоков от имени других пользователей.
  Необходимо если Spark работает под другим UNIX пользователем, чем DataNode.
-->
<property>
  <name>dfs.block.local-path-access.user</name>
  <value>spark,yarn,mapred</value>
</property>

<!--
  Выполнять checksum проверку при short-circuit чтении.
  false: пропустить проверку для максимальной скорости (риск corrupt data)
  true: выполнять проверку (рекомендуется для production)
-->
<property>
  <name>dfs.client.read.shortcircuit.skip.checksum</name>
  <value>false</value>
</property>

Проброс в SparkSession (если hdfs-site.xml уже содержит настройки, явная конфигурация не нужна):

spark = SparkSession.builder \
    .config("spark.hadoop.dfs.client.read.shortcircuit", "true") \
    .config("spark.hadoop.dfs.domain.socket.path", "/var/lib/hadoop-hdfs/dn_socket") \
    .getOrCreate()

Проверка работы Short-Circuit Reads

Убедиться, что Short-Circuit Reads реально используется, можно через JMX-метрики DataNode:

# Метрики DataNode через JMX (порт 9864 для HDFS 3.x)
curl -s "http://datanode1:9864/jmx?qry=Hadoop:service=DataNode,name=DataNodeActivity" \
  | python3 -m json.tool \
  | grep -E "BlocksRead|ShortCircuit"

# Ожидаемый вывод при активном Short-Circuit:
# "BlocksRead" : 15234,
# "ShortCircuitLocalReadsRead" : 14891,  # ~97% от BlocksRead
# "ShortCircuitLocalBytesRead" : 1234567890

Если ShortCircuitLocalReadsRead близко к BlocksRead - Short-Circuit работает. Если близко к нулю - проверьте доступность Unix Domain Socket и права пользователей.


5. Риски деградации локальности: Data Skew и Speculative Execution

Идеальный мир, где каждый блок равномерно распределён по кластеру и каждый Executor всегда получает NODE_LOCAL задачи, редко встречается на практике.

Проблема неравномерного распределения блоков (Data Hotspots)

Представьте сценарий: аналитическая таблица orders за 3 года, партиционированная по месяцам. 90% запросов читают данные за последние 6 месяцев. Если DataNode1 и DataNode2 случайно оказались первичными носителями большинства блоков этих партиций - они становятся «горячими» (Hotspot).

При большом spark.locality.wait в ситуации Hotspot возникает парадоксальная проблема: Spark ждёт освобождения ресурсов на перегруженных DataNode1/2, при этом DataNode3/4 с их Executor'ами простаивают. Итоговое время выполнения Stage растёт.

Диагностика Hotspot через hdfs dfsadmin:

# Просмотр распределения данных по DataNode
hdfs dfsadmin -report

# Пример вывода с явным дисбалансом:
# DataNode1: DFS Used%: 78%  ← Hotspot!
# DataNode2: DFS Used%: 71%
# DataNode3: DFS Used%: 12%
# DataNode4: DFS Used%: 15%

# Запуск балансировщика HDFS
hdfs balancer -threshold 10
# threshold 10: балансировка завершена если разница < 10%

Взаимодействие с Speculative Execution

Speculative Execution - механизм Spark, при котором задачи, работающие значительно медленнее остальных в Stage, перезапускаются параллельно на других Executor'ах.

spark = SparkSession.builder \
    # Включить спекулятивное выполнение
    .config("spark.speculation", "true") \

    # Задача считается «отстающей», если выполняется в multiplier
    # раз дольше медианы всех задач в Stage.
    # Дефолт: 1.5 (на 50% медленнее медианы = кандидат)
    .config("spark.speculation.multiplier", "1.5") \

    # Запускать спекуляцию только если завершено не менее fraction задач Stage.
    # При 0.75: ждём завершения 75% задач прежде чем решить «задача медленная».
    .config("spark.speculation.quantile", "0.75") \

    .getOrCreate()

Конфликт Speculative Execution с Delay Scheduling: спекулятивная копия задачи запускается немедленно на любом доступном Executor'е (уровень ANY), чтобы выполниться быстрее. Это противоречит Delay Scheduling, который ждёт NODE_LOCAL слота.

Правило: при включённом Speculation увеличивайте spark.speculation.quantile до 0.9+ чтобы не запускать спекуляцию по случайным замедлениям, связанным с ожиданием локальности. Иначе медленные задачи (которые медленные именно потому что ждут NODE_LOCAL слот) будут немедленно перезапущены на ANY, отменяя всю работу Delay Scheduling.

Влияние Shuffle на локальность: почему она теряется

Data Locality - это свойство входных данных Stage 0 (чтение с HDFS). После первого Shuffle она полностью теряется.

После Shuffle каждая задача в следующем Stage читает данные из Shuffle партиций, которые разбросаны по всему кластеру. Никакого «предпочтительного расположения» для такой задачи нет.

Как вернуть PROCESS_LOCAL через кеш: если результат тяжёлого агрегирующего Stage нужен несколько раз, кешируйте его:

# Первое выполнение: Stage 0 (NODE_LOCAL) + Stage 1 (ANY после Shuffle)
silver_df = (
    spark.read.parquet("hdfs://cluster/bronze/events/")
    .filter(df.event_type == "purchase")
    .groupBy("user_id")
    .agg({"amount": "sum", "event_id": "count"})
)

# Кешируем - данные хранятся в памяти Executor'ов
silver_df.cache()
silver_df.count()  # материализуем кеш

# Последующие операции с silver_df получат PROCESS_LOCAL:
# данные уже в heap Executor'ов, никакого I/O нет
result1 = silver_df.join(users_df, "user_id")
result2 = silver_df.join(products_df, "user_id")

6. Мониторинг и аудит локальности через Spark UI и метрики

Теоретическое понимание уровней локальности бесполезно без инструментов, позволяющих видеть, что реально происходит в кластере.

Инспекция вкладки Stages в Spark UI

В Spark UI (по умолчанию http://driver-host:4040) на вкладке Stages для каждого Stage доступна детальная таблица задач. Ключевые столбцы для анализа локальности:

Столбец Что показывает
Locality Level Уровень локальности конкретной задачи
Scheduler Delay Время от готовности задачи до её старта (включает ожидание локальности)
Getting Result Time Время передачи результата задачи драйверу
Input Size Объём прочитанных данных
Duration Реальное время выполнения задачи

Паттерны для анализа:

Здоровый профиль локальности:

Locality Level Distribution:
  NODE_LOCAL: 847 задач (92%)
  RACK_LOCAL:  62 задачи  (7%)
  ANY:         11 задач   (1%)
Среднее Scheduler Delay: 45ms

Нездоровый профиль - кластер гоняет данные по сети:

Locality Level Distribution:
  NODE_LOCAL:  123 задачи (13%)
  RACK_LOCAL:  234 задачи (25%)
  ANY:         623 задачи (66%)  ← проблема!
Среднее Scheduler Delay: 3200ms  ← долгое ожидание

Программный доступ к статистике локальности через Spark History Server API:

import requests
from collections import Counter


def get_locality_distribution(app_id: str, stage_id: int,
                               history_server_url: str) -> dict:
    """
    Получает распределение уровней локальности для заданного Stage.
    Полезно для автоматизированного аудита после завершения джобы.
    """
    url = (f"{history_server_url}/api/v1/applications/{app_id}"
           f"/stages/{stage_id}/0/taskList")
    tasks = requests.get(url).json()

    locality_counts = Counter(task["taskLocality"] for task in tasks)
    total = sum(locality_counts.values())

    distribution = {
        level: {
            "count": count,
            "percentage": round(count / total * 100, 1)
        }
        for level, count in locality_counts.items()
    }

    # Задачи с аномально высоким Scheduler Delay
    high_delay_threshold_ms = 5000
    slow_tasks = [
        {
            "taskId": t["taskId"],
            "locality": t["taskLocality"],
            "schedulerDelayMs": t["taskMetrics"]["schedulerDelay"],
            "durationMs": t["duration"],
        }
        for t in tasks
        if t.get("taskMetrics", {}).get("schedulerDelay", 0) > high_delay_threshold_ms
    ]

    return {
        "stageId": stage_id,
        "totalTasks": total,
        "localityDistribution": distribution,
        "slowTasksCount": len(slow_tasks),
        "slowTasks": slow_tasks[:5],
    }

Диагностика через логи планировщика

Логи Spark Driver на уровне DEBUG для org.apache.spark.scheduler.TaskSetManager содержат детальную информацию о каждом решении по локальности:

# Запуск Spark с детальным логированием планировщика
spark-submit \
  --conf "spark.driver.extraJavaOptions=-Dlog4j.logger.org.apache.spark.scheduler=DEBUG" \
  your_job.py

Характерные строки в логах:

# Задача получила NODE_LOCAL:
INFO  TaskSetManager: Starting task 5.0 in stage 0.0 (TID 15,
  datanode3.example.com, executor 3, partition 5, NODE_LOCAL, 5242 bytes)

# Снижение требований к локальности:
WARN  TaskSetManager: No NODE_LOCAL tasks are launchable, waiting 3000ms
INFO  TaskSetManager: Level for task 4 has been downgraded from NODE_LOCAL
  to RACK_LOCAL because it waited 3000ms without finding a suitable executor

# Спекулятивный запуск задачи:
INFO  TaskSetManager: Marking task 12 in stage 1.0 (on server2) as
  speculative because it ran more than 15.0s
INFO  TaskSetManager: Starting task 12.1 in stage 1.0 (TID 156,
  server5.example.com, executor 8, partition 12, ANY, 5102 bytes)

Ключевая метрика: Fetch Wait Time

Fetch Wait Time - время, которое Spark-задача тратит на ожидание получения Shuffle-данных от других Executor'ов. Это метрика сетевого bottleneck.

Если Fetch Wait Time составляет 40–60% от Duration задачи - кластерная сеть или DataNode диски являются узким местом. Это часто происходит при уровне ANY, когда задача получает данные от десятков Executor'ов по сети.

# Пороги деградации для автоматического мониторинга
LOCALITY_THRESHOLDS = {
    "node_local_min_pct": 70,        # минимум 70% задач NODE_LOCAL
    "any_max_pct": 15,               # не более 15% задач ANY
    "fetch_wait_max_ratio": 0.3,     # Fetch Wait Time < 30% от Duration
    "scheduler_delay_max_ms": 3000,  # Scheduler Delay < 3 сек
}

7. Практика: тюнинг локальности и лабораторная работа

Конфигурация Spark для максимальной HDFS локальности

# production_spark_config.py
# Производственная конфигурация SparkSession с учётом Data Locality

from pyspark.sql import SparkSession


def create_spark_hdfs_optimized(
    app_name: str,
    nameservice: str = "mycluster",
    task_duration_seconds: float = 30.0,
    enable_short_circuit: bool = True,
    enable_speculation: bool = False,
) -> SparkSession:
    """
    Создаёт SparkSession с настройками для максимальной HDFS локальности.

    task_duration_seconds: ожидаемое среднее время выполнения одной задачи.
    Используется для расчёта оптимального spark.locality.wait.

    Правило выбора locality.wait:
    - Если задача < 5 сек:  wait = 0 (ожидание бессмысленно)
    - Если задача 5-30 сек: wait = 3-5 сек
    - Если задача > 30 сек: wait = 5-15 сек
    """
    if task_duration_seconds < 5:
        locality_wait = "0"
        locality_wait_node = "0"
    elif task_duration_seconds < 30:
        locality_wait = "3s"
        locality_wait_node = "5s"
    else:
        locality_wait = "5s"
        locality_wait_node = "15s"

    builder = SparkSession.builder \
        .appName(app_name) \
        .config("spark.hadoop.fs.defaultFS", f"hdfs://{nameservice}") \
        # PROCESS_LOCAL при первом чтении HDFS недостижим → ставим 0
        .config("spark.locality.wait.process", "0") \
        .config("spark.locality.wait.node", locality_wait_node) \
        .config("spark.locality.wait.rack", locality_wait) \
        .config("spark.locality.wait", locality_wait) \
        .config("spark.task.maxFailures", "4")

    if enable_short_circuit:
        builder = builder \
            .config("spark.hadoop.dfs.client.read.shortcircuit", "true") \
            .config("spark.hadoop.dfs.domain.socket.path",
                    "/var/lib/hadoop-hdfs/dn_socket") \
            .config("spark.hadoop.dfs.client.read.shortcircuit.skip.checksum",
                    "false")

    if enable_speculation:
        builder = builder \
            .config("spark.speculation", "true") \
            .config("spark.speculation.multiplier", "2.0") \
            .config("spark.speculation.quantile", "0.9")
    else:
        builder = builder.config("spark.speculation", "false")

    return builder.getOrCreate()


# Для batch ETL с долгими задачами (агрегации, joins на десятках GB)
spark_batch = create_spark_hdfs_optimized(
    app_name="daily-etl",
    task_duration_seconds=60.0,
    enable_short_circuit=True,
)

# Для Structured Streaming (критично: wait=0 чтобы не задерживать микробатчи)
spark_streaming = create_spark_hdfs_optimized(
    app_name="kafka-to-hdfs-stream",
    task_duration_seconds=1.0,
    enable_short_circuit=True,
    enable_speculation=False,  # Speculation и Streaming несовместимы
)

Лабораторная работа: сравнение locality.wait = 0 vs 10 секунд

# lab_locality_wait_comparison.py
# Сравниваем влияние spark.locality.wait на производительность и
# распределение уровней локальности.

import time
from dataclasses import dataclass
from pyspark.sql import SparkSession
from pyspark.sql import functions as F


@dataclass
class BenchmarkResult:
    config_name: str
    locality_wait: str
    duration_seconds: float


def run_benchmark(spark: SparkSession, data_path: str,
                  config_name: str, locality_wait: str) -> BenchmarkResult:
    """
    Запускает стандартный аналитический запрос с заданным locality.wait.
    Измеряет время выполнения.
    """
    spark.conf.set("spark.locality.wait", locality_wait)
    spark.conf.set("spark.locality.wait.node", locality_wait)

    print(f"\nКонфигурация: {config_name}, locality.wait = {locality_wait}")
    start_time = time.time()

    events = spark.read.parquet(f"{data_path}/events/")
    users = spark.read.parquet(f"{data_path}/users/")

    result = (
        events
        .filter(F.col("event_type").isin("purchase", "checkout"))
        .join(users.select("user_id", "country", "segment"), "user_id", "left")
        .groupBy("country", "segment",
                 F.date_trunc("day", "event_time").alias("day"))
        .agg(
            F.count("*").alias("event_count"),
            F.sum("amount").alias("total_revenue"),
            F.countDistinct("session_id").alias("unique_sessions"),
        )
        .orderBy("day", "country")
    )

    row_count = result.count()
    duration = time.time() - start_time

    print(f"Результат: {row_count} строк за {duration:.1f} секунд")
    return BenchmarkResult(config_name, locality_wait, duration)


def generate_benchmark_data(spark: SparkSession, output_path: str,
                             n_events: int = 5_000_000) -> None:
    """Генерирует тестовый датасет для бенчмарка на HDFS."""
    print(f"Генерируем {n_events:,} событий → {output_path}")

    events = spark.range(n_events).select(
        F.col("id").alias("event_id"),
        (F.rand() * 100_000).cast("long").alias("user_id"),
        F.array(
            F.lit("purchase"), F.lit("view"), F.lit("checkout"), F.lit("search")
        ).getItem((F.rand() * 4).cast("int")).alias("event_type"),
        (F.rand() * 1000).alias("amount"),
        F.expr("uuid()").alias("session_id"),
        F.current_timestamp().alias("event_time"),
    )

    # Создаём много небольших файлов для теста локальности
    events.repartition(200).write.mode("overwrite").parquet(f"{output_path}/events/")

    users = spark.range(100_000).select(
        F.col("id").alias("user_id"),
        F.array(F.lit("RU"), F.lit("DE"), F.lit("US")).getItem(
            (F.rand() * 3).cast("int")
        ).alias("country"),
        F.array(F.lit("premium"), F.lit("standard"), F.lit("free")).getItem(
            (F.rand() * 3).cast("int")
        ).alias("segment"),
    )
    users.coalesce(10).write.mode("overwrite").parquet(f"{output_path}/users/")
    print("Данные записаны")


def print_comparison(results: list[BenchmarkResult]) -> None:
    print("\n" + "=" * 60)
    print("СРАВНЕНИЕ КОНФИГУРАЦИЙ LOCALITY.WAIT")
    print("=" * 60)
    print(f"{'Конфигурация':<30} {'Wait':>6} {'Время':>10}")
    print("-" * 60)
    for r in results:
        print(f"{r.config_name:<30} {r.locality_wait:>6} {r.duration_seconds:>9.1f}s")

    fastest = min(results, key=lambda r: r.duration_seconds)
    print(f"\nБыстрейшая конфигурация: {fastest.config_name}")


if __name__ == "__main__":
    spark = SparkSession.builder \
        .master("yarn") \
        .appName("locality-wait-benchmark") \
        .config("spark.hadoop.fs.defaultFS", "hdfs://mycluster") \
        .config("spark.executor.instances", "20") \
        .config("spark.executor.cores", "4") \
        .config("spark.executor.memory", "8g") \
        .config("spark.hadoop.dfs.client.read.shortcircuit", "true") \
        .config("spark.hadoop.dfs.domain.socket.path",
                "/var/lib/hadoop-hdfs/dn_socket") \
        .getOrCreate()

    DATA_PATH = "hdfs://mycluster/lab/locality-benchmark/"
    generate_benchmark_data(spark, DATA_PATH, n_events=5_000_000)

    configs = [
        ("Без ожидания (wait=0)", "0"),
        ("Базовый (wait=3s)", "3s"),
        ("Агрессивный (wait=10s)", "10s"),
        ("Максимальный (wait=30s)", "30s"),
    ]

    results = []
    for name, wait in configs:
        result = run_benchmark(spark, DATA_PATH, name, wait)
        results.append(result)
        time.sleep(5)

    print_comparison(results)
    spark.stop()

Диагностическое дерево: что делать при плохой локальности

Команды для быстрой диагностики состояния кластера

# 1. Распределение данных по DataNode (ищем Hotspot)
hdfs dfsadmin -report | grep -A5 "DataNode"

# 2. Балансировка данных (при дисбалансе > 15%)
hdfs balancer -threshold 10 &
tail -f /var/log/hadoop/hdfs/hadoop-hdfs-balancer-*.log

# 3. Block locations конкретного файла
hdfs fsck /data/events/part-00001.parquet -files -blocks -locations

# 4. Проверить Short-Circuit метрики DataNode
curl -s "http://datanode1.example.com:9864/jmx" \
  | python3 -c "
import sys, json
data = json.load(sys.stdin)
for bean in data['beans']:
    if 'DataNodeActivity' in bean.get('name', ''):
        for k in ['BlocksRead', 'ShortCircuitLocalReadsRead']:
            print(f'{k}: {bean.get(k, 0):,}')
"

# 5. Topology: rack ID каждого DataNode
hdfs dfsadmin -printTopology

# 6. Работающие Spark приложения
yarn application -list -appStates RUNNING

Итоговая матрица настроек для разных сценариев

Сценарий locality.wait.node locality.wait.rack short.circuit speculation
Batch ETL, долгие задачи (> 60 сек) 10–15 сек 5 сек true false
Интерактивная аналитика, короткие задачи 0–1 сек 0 true false
Structured Streaming, микробатчи 0 0 true false
Кластер перегружен (CPU > 80%) 0–2 сек 0 true true
Облачный кластер (S3/MinIO) 0 0 false false
On-Premise, свежий кластер (< 30% CPU) 5–10 сек 3 сек true false

Итоги: ключевые выводы о Data Locality

Data Locality - это не просто технический трюк, а фундаментальный принцип, от которого зависит эффективность всего HDFS-кластера.

Иерархия уровней: PROCESS_LOCAL (кеш в JVM) → NODE_LOCAL (локальный диск) → RACK_LOCAL (та же стойка) → ANY (по сети). Разница в производительности между крайними уровнями - десятки раз на больших объёмах.

Delay Scheduling - компромисс между скоростью запуска и качеством локальности. Ключевой параметр spark.locality.wait.node: увеличивайте для долгих задач на нагруженных кластерах, устанавливайте в 0 для стриминга и коротких задач.

Short-Circuit Reads - обязательная оптимизация для HDFS-кластеров. Устраняет TCP overhead при NODE_LOCAL чтении, даёт 10–30% прироста производительности при минимальных изменениях конфигурации.

Shuffle разрушает HDFS-локальность для последующих Stage'ов. Кеш (.cache()) - способ вернуть PROCESS_LOCAL для данных, которые читаются несколько раз после агрегации.

В облачных и S3-based Lakehouse архитектурах (S3, MinIO, Ceph, Ozone) классическая HDFS-локальность не применима - Compute и Storage разделены по определению. Устанавливайте spark.locality.wait = 0 и фокусируйтесь на других оптимизациях: column pruning, predicate pushdown, partition pruning, кеширование.