Data Locality: перемести код к данным, а не данные к коду
Уровни локальности (PROCESS_LOCAL → ANY), Delay Scheduling, диагностика в Spark UI и тюнинг locality.wait для HDFS и объектных хранилищ.
Философия: почему данные должны ждать код, а не наоборот¶
В традиционных системах данные перемещаются к вычислению: запрос идёт на сервер, который хранит таблицу, или данные копируются на машину с аналитическим ПО. Для гигабайт это нормально.
В распределённых системах, где таблица занимает терабайты, это невозможно. Стоимость передачи данных по сети часто превышает стоимость самих вычислений:
Иерархия задержек на практике:
| Операция | Пропускная способность | Время для 1 GB |
|---|---|---|
| L3 cache → CPU | ~200 GB/s | ~5 мс |
| DRAM | ~50–100 GB/s | 10–20 мс |
| Локальный NVMe | ~5 GB/s | ~200 мс |
| Сеть внутри стойки (25 GbE) | ~3 GB/s | ~340 мс |
| Сеть между стойками (10 GbE) | ~1 GB/s | ~1 с |
| S3 / объектное хранилище | ~0.1–0.5 GB/s | 2–10 с |
Вывод: передача 10 GB через межстоечную сеть занимает 10–100 секунд. Передача Python UDF-кода (10–100 KB) - миллисекунды. Принцип Data Locality - отправить код к данным, а не наоборот.
Архитектура хранения данных: где живут блоки¶
Чтобы понять локальность, нужно понять, как данные хранятся физически.
HDFS: блоки с репликацией¶
Каждый HDFS-блок (по умолчанию 128 MB) реплицируется на 3 узла (replication factor). Это означает: для любого блока есть 3 узла, на которых Task может работать без сетевой передачи данных.
Объектное хранилище: нет физической локальности¶
S3, MinIO, Ceph RGW - это отдельные сервисы. Данные доступны через HTTP-API с любого узла. Понятие "блок на узле" отсутствует.
Для S3 все Executor'ы одинаково далеко от данных. Spark устанавливает уровень NO_PREF.
Уровни Data Locality¶
Spark определяет пять уровней локальности (от лучшего к худшему):
PROCESS_LOCAL: идеальный сценарий¶
Task запускается в том же JVM-процессе (Executor), в чьей памяти уже лежит партиция данных. Это происходит при:
df.cache()/df.persist()- данные закэшированы вStorageMemoryExecutor'а- Повторное использование broadcast переменных (закэшированы локально)
# Первый Action - данные читаются из HDFS/S3, кэшируются в памяти Executor'а
df = spark.read.parquet("/data/events").cache()
df.count() # NODE_LOCAL (чтение с локального диска HDFS DataNode)
# Второй Action - данные уже в памяти Executor'а
df.groupBy("region").count() # PROCESS_LOCAL
df.filter("date >= '2024-01-01'").show() # PROCESS_LOCAL
NODE_LOCAL: локальный диск¶
Task запускается на том же физическом узле, что и DataNode, хранящий нужный блок. Данные читаются с локального диска - без сети. В HDFS-кластерах это самый частый уровень для первого чтения.
Условие: Executor и DataNode сопримещены (co-located) на одном хосте. В YARN это настраивается автоматически - NodeManager знает, какие HDFS-блоки есть на этом хосте.
RACK_LOCAL: внутристоечная сеть¶
Нужный блок находится на другом узле, но в той же физической серверной стойке. Трафик проходит через Top-of-Rack (ToR) коммутатор, но не выходит выше. Обычно 10–25 Gbps.
Spark учитывает топологию сети через spark.network.topology.script - скрипт, возвращающий rack-имя для каждого IP-адреса.
ANY: самый медленный¶
Task запускается на любом свободном Executor'е без учёта расположения данных. Данные передаются по межстоечной или даже междатацентровой сети.
Когда это происходит неизбежно:
- Все узлы с нужными данными перегружены и свободных слотов нет
- Data Skew: одна партиция огромная, остальные Executor'ы уже закончили
- S3 / объектное хранилище:
NO_PREF→ планировщик не ждёт и сразу назначает на любой свободный
Delay Scheduling: алгоритм ожидания¶
Когда предпочтительный узел занят, Spark не сдаётся сразу. Он ждёт - вдруг освободится нужный Executor. Этот механизм называется Delay Scheduling.
Параметры Delay Scheduling¶
# Общее время ожидания (применяется к каждому уровню, если не переопределено)
spark.conf.set("spark.locality.wait", "3s") # default: 3s
# Тонкая настройка по уровням:
spark.conf.set("spark.locality.wait.process", "3s") # ждать PROCESS_LOCAL
spark.conf.set("spark.locality.wait.node", "3s") # ждать NODE_LOCAL
spark.conf.set("spark.locality.wait.rack", "3s") # ждать RACK_LOCAL
Логика последовательного снижения уровня:
Трейдофф: ждать или не ждать¶
Увеличить locality.wait имеет смысл когда:
- Задачи тяжёлые (минуты выполнения) - ожидание 10–15 с незначительно
- Кластер временно перегружен (burst load) - нужные узлы скоро освободятся
- Данные на HDFS с репликацией 3, Executor'ы часто сопримещены с DataNode
Уменьшить locality.wait (вплоть до 0) имеет смысл когда:
- Данные в S3 / объектном хранилище - ждать бессмысленно (
NO_PREF) - Задачи короткие (секунды) - время ожидания сопоставимо со временем выполнения
- Кластер сильно загружен, нужные Executor'ы никогда не освобождаются
# S3-кластер: не ждём вообще
spark.conf.set("spark.locality.wait", "0s")
# HDFS-кластер с тяжёлыми задачами: ждём дольше
spark.conf.set("spark.locality.wait", "10s")
Shuffle разрушает локальность¶
Shuffle - момент, когда данные принудительно перераспределяются между Executor'ами. Любая локальность, достигнутая при чтении исходных данных, теряется.
Важный вывод: groupBy, join, repartition, distinct - всё это shuffle. После shuffle Spark начинает с нуля: Reduce-фаза работает со смешанными данными со всех узлов, уровень ANY.
Narrow vs Wide Transformations¶
# Narrow: нет shuffle, локальность сохраняется
df.filter(col("status") == "active") # каждая партиция → один Task на том же узле
df.withColumn("upper", upper(col("name"))) # то же самое
# Wide: shuffle, локальность теряется
df.groupBy("region").count() # map → shuffle → reduce, ALL локальность сброшена
df.join(other, "user_id") # если не broadcast - shuffle обоих датасетов
Broadcast Join: оптимизация через отказ от shuffle¶
Broadcast Join - способ сохранить локальность для join-операций с маленькой таблицей.
Маленькая таблица копируется на каждый Executor. Большая таблица остаётся на своих местах. Никакого shuffle для большой таблицы - PROCESS_LOCAL или NODE_LOCAL.
from pyspark.sql.functions import broadcast
# Явный broadcast hint
result = large_df.join(broadcast(small_df), "product_id")
# Автоматический (если small_df < spark.sql.autoBroadcastJoinThreshold = 10 MB)
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "50m")
Объектное хранилище: другой мир¶
S3, MinIO, Azure Blob, GCS - это не HDFS. У них нет Data Locality. Все Executor'ы одинаково удалены от данных.
Локальный кэш поверх S3¶
Для повторяющихся запросов к одним и тем же данным можно использовать Alluxio или локальный SSD-кэш на Executor'ах. После первого чтения данные кэшируются на узле, и повторные запросы получают NODE_LOCAL вместо NO_PREF.
# Spark + Alluxio: прозрачный кэш поверх S3
df = spark.read.parquet("alluxio://master:19998/data/events/")
# Первый раз: читает из S3, кэширует на Alluxio-узлах
# Повторно: NODE_LOCAL с Alluxio worker'а
Факторы, убивающие локальность¶
1. Нехватка ресурсов на нужных узлах¶
Самая частая причина деградации до ANY. Если все ядра на Node1 заняты другими задачами, и locality.wait истёк - задача уходит на свободный Node5.
2. Data Skew: одна огромная партиция¶
Если одна партиция содержит 90% данных, задача для неё выполняется долго. Все остальные Executor'ы простаивают или берут другие задачи, не связанные с этими данными.
3. Слишком много мелких партиций¶
При 10 000 партиций планировщик делает 10 000 запросов о локальности. Накладные расходы на планирование начинают перевешивать выгоду от локальности.
4. Dynamic Allocation и масштабирование¶
Dynamic Allocation добавляет новые Executor'ы, когда есть задачи в очереди. Новые Executor'ы не знают о расположении данных - они получают уровень ANY до тех пор, пока не закэшируют что-либо локально.
# Dynamic Allocation снижает locality при агрессивном масштабировании
spark.conf.set("spark.dynamicAllocation.enabled", "true")
spark.conf.set("spark.dynamicAllocation.minExecutors", "4") # базовый уровень знает locality
spark.conf.set("spark.dynamicAllocation.maxExecutors", "20") # новые - без locality
Диагностика в Spark UI¶
Вкладка Stages → Tasks¶
Что смотреть при ANY:
- Shuffle Read Size (вкладка Tasks) - сколько байт передано по сети для этой задачи
- Task Duration - аномально долгие задачи при
ANYуказывают на сетевое узкое место - Input Size - если
ANYи Input Size большой, это прямые потери на сетевую передачу
Метрики в коде¶
# После выполнения job смотреть статус через SparkContext
sc = spark.sparkContext
status = sc.statusTracker()
for stage_id in status.getActiveStageIds():
info = status.getStageInfo(stage_id)
print(f"Stage {stage_id}: {info.numActiveTasks()} tasks")
# Или через программный доступ к истории:
# spark.sparkContext._jsc.sc().statusTracker().getJobIdsForGroup(...)
Cache/Persist и локальность¶
Кэширование - главный инструмент для достижения PROCESS_LOCAL:
from pyspark import StorageLevel
# MEMORY_AND_DISK: PROCESS_LOCAL, при нехватке памяти - NODE_LOCAL
df.persist(StorageLevel.MEMORY_AND_DISK)
# MEMORY_ONLY: PROCESS_LOCAL; если не помещается - задача перечитывает данные
df.persist(StorageLevel.MEMORY_ONLY)
# DISK_ONLY: NODE_LOCAL (данные на локальном диске Executor'а)
df.persist(StorageLevel.DISK_ONLY)
Важно: кэш df.cache() прибит к конкретным Executor'ам. Если Executor умрёт (или уйдёт при Dynamic Allocation), данные потеряются - блок будет перечитан и снова закэширован, но уже на другом Executor'е.
# Антипаттерн: кэшировать данные, которые используются один раз
df.cache()
df.count() # кэш заполнен
result = df.groupBy("region").count() # используем кэш
df.unpersist() # ВАЖНО: освободить память после использования
Тюнинг locality.wait: практические рекомендации¶
Антипаттерны, разрушающие локальность¶
# Антипаттерн 1: лишний repartition в начале pipeline
df = spark.read.parquet("/hdfs/data/") # NODE_LOCAL
df = df.repartition(200) # shuffle → ANY, данные перемешаны
df.filter(...) # теперь данные не там, где начинались
# Лучше: coalesce или увеличить начальные партиции
df = spark.read.parquet("/hdfs/data/")
df = df.coalesce(200) # без shuffle, только уменьшение числа партиций
# Антипаттерн 2: join без broadcast при маленькой таблице
# large_df.join(small_df, ...) → shuffle обоих → ANY
# Лучше:
large_df.join(broadcast(small_df), ...) # только large_df остаётся на месте
# Антипаттерн 3: бесконечно наращивать партиции
df.repartition(10000) # 10000 задач планировщик запускает с 10000 locality-запросов
# Overhead планировщика перевешивает выгоду
Практика: диагностика уровня локальности¶
1. Посмотреть locality level через explain и UI¶
# Создаём тестовый датасет на HDFS (или имитируем)
df = spark.read.parquet("/data/large_table/")
df.cache()
df.count() # первое чтение: NODE_LOCAL (HDFS)
# Второе чтение
df.groupBy("region").count().show()
# В Spark UI → Stages → Tasks: смотрим Locality Level
# Ожидаем PROCESS_LOCAL для задач на закэшированных партициях
2. Сравнить время с кэшем и без¶
import time
# Без кэша (NODE_LOCAL от HDFS или NO_PREF от S3)
t0 = time.time()
spark.read.parquet("/data/events/").filter("year=2024").count()
t_cold = time.time() - t0
# С кэшем (PROCESS_LOCAL)
df = spark.read.parquet("/data/events/").filter("year=2024").cache()
df.count() # прогрев
t0 = time.time()
df.count() # из кэша
t_cached = time.time() - t0
print(f"Без кэша: {t_cold:.1f}s")
print(f"С кэшем: {t_cached:.1f}s")
print(f"Ускорение: {t_cold/t_cached:.1f}×")
3. Анализ locality через Spark History Server¶
# После завершения задания locality level виден в REST API History Server
import requests
app_id = "application_1234567890_0001"
url = f"http://spark-history:18080/api/v1/applications/{app_id}/stages"
stages = requests.get(url).json()
for stage in stages:
tasks_url = f"http://spark-history:18080/api/v1/applications/{app_id}/stages/{stage['stageId']}/0/taskList"
tasks = requests.get(tasks_url).json()
levels = {}
for task in tasks:
lvl = task.get("taskLocality", "UNKNOWN")
levels[lvl] = levels.get(lvl, 0) + 1
print(f"Stage {stage['stageId']}: {levels}")
Лучшие практики¶
| Ситуация | Рекомендация |
|---|---|
| HDFS-кластер | Держать locality.wait = 3s (дефолт), сопримещать DataNode + NodeManager |
| S3 / объектное хранилище | locality.wait = 0, фокус на partition pruning и columnar форматах |
| Повторяющиеся запросы к одним данным | df.cache() - переключает на PROCESS_LOCAL |
| Join с маленькой таблицей | Всегда broadcast(small_df) - избегает shuffle большой таблицы |
| Тяжёлые задачи (> 1 мин) | Увеличить locality.wait до 10–15s |
| Короткие задачи (< 5c) | Уменьшить locality.wait до 1s или 0 |
| Dynamic Allocation | Держать minExecutors с запасом для горячих данных |
| Data Skew | Salting или AQE Skew Join - иначе locality бесполезна для skewed партиций |
Резюме¶
Data Locality - способность Spark запускать Task рядом с данными, а не тащить данные к Task. Иерархия от лучшего к худшему:
PROCESS_LOCAL (данные в памяти того же Executor) → NODE_LOCAL (на том же узле) → NO_PREF (нет предпочтения, S3) → RACK_LOCAL (та же стойка) → ANY (любой узел, полная сеть).
Delay Scheduling даёт планировщику время подождать нужный Executor. При S3 ожидание бессмысленно - locality.wait = 0. При HDFS с тяжёлыми задачами - стоит подождать.
Shuffle необратимо разрушает локальность. Broadcast Join, кэширование и partition pruning - главные инструменты её сохранения.