Data Locality: перемести код к данным, а не данные к коду

Уровни локальности (PROCESS_LOCAL → ANY), Delay Scheduling, диагностика в Spark UI и тюнинг locality.wait для HDFS и объектных хранилищ.

core optimization

Философия: почему данные должны ждать код, а не наоборот

В традиционных системах данные перемещаются к вычислению: запрос идёт на сервер, который хранит таблицу, или данные копируются на машину с аналитическим ПО. Для гигабайт это нормально.

В распределённых системах, где таблица занимает терабайты, это невозможно. Стоимость передачи данных по сети часто превышает стоимость самих вычислений:

Иерархия задержек на практике:

Операция Пропускная способность Время для 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() - данные закэшированы в StorageMemory Executor'а
  • Повторное использование 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 - главные инструменты её сохранения.