Block Locality: как Spark scheduler использует data-locality при планировании
Фундаментальный разбор Data Locality в HDFS-кластерах: физика блочного хранения, иерархия уровней PROCESS_LOCAL → ANY, алгоритм Delay Scheduling, механизм Short-Circuit Reads, влияние Shuffle и кеширования, мониторинг через Spark UI и тюнинг под SLA.
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):
- Ждём
spark.locality.wait.processсекунд для PROCESS_LOCAL → если нет, пробуем NODE_LOCAL - Ждём
spark.locality.wait.nodeсекунд для NODE_LOCAL → если нет, пробуем RACK_LOCAL - Ждём
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-стек:
- DataNode читает блок с диска в свой Java heap
- DataNode сериализует данные в TCP-буфер ядра ОС
- Spark Executor читает из TCP-буфера в свой Java heap
- Данные оказываются скопированными дважды: 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, кеширование.