Spark vs Hadoop MapReduce: in-memory вычисления и Lazy Evaluation

Почему Spark в 10–100 раз быстрее MapReduce, как работают ленивые вычисления и как Catalyst оптимизирует запросы до начала выполнения.

core optimization

В предыдущем уроке мы разобрали анатомию кластера: Driver строит DAG, Task Scheduler распределяет задачи по Executor-ам, данные живут в Partition-ах. Теперь разберём почему эта архитектура в 10–100 раз быстрее MapReduce и как механизм Lazy Evaluation позволяет Spark оптимизировать план ещё до первой строки вычислений.


Наследие MapReduce: диск как узкое горлышко

MapReduce появился в Google в 2004 году как ответ на вопрос «как обработать петабайт данных на тысячах ненадёжных машин?». Решение элегантное: разбей задачу на Map (обработай каждую запись независимо) и Reduce (агрегируй результаты по ключу), запускай везде параллельно. Apache Hadoop реализовал эту идею как open-source, и к 2010 году стал стандартом для Big Data.

Проблема не в самой модели - проблема в обязательном дисковом I/O.

Каждая MapReduce-задача обязана записать промежуточные данные на HDFS перед тем, как передать их следующему этапу. Это не баг - это гарантия отказоустойчивости: если узел упал посреди Reduce, следующий Job читает уже записанный результат, а не теряет его. За надёжность платят скоростью.

Для одного SQL-запроса с двумя JOIN-ами Hive компилировал в три отдельных MapReduce Job-а - каждый со своим чтением HDFS, записью промежуточных файлов и повторным чтением. Суммарно: 6 операций с распределённым диском на один запрос.

Итеративные алгоритмы: катастрофа

Машинное обучение работает итеративно: алгоритм многократно проходит по данным, корректируя веса. Логистическая регрессия - 100 итераций, K-Means - 50 итераций, PageRank - 30 итераций.

В MapReduce каждая итерация = отдельный Job = чтение HDFS + запись HDFS. 100 итераций = 200 операций с HDFS на холодном диске.

Бенчмарк из оригинальной статьи о Spark (Zaharia et al., Berkeley AMPLab, 2012):

Задача Hadoop MapReduce Apache Spark
Логистическая регрессия (100 итераций) ~110 сек / итерация ~0.9 сек / итерация
K-Means кластеризация ~26 мин ~80 сек

Разница не в качестве кода - разница в том, что Spark держит данные в памяти между итерациями.


Spark: промежуточные данные в памяти

Spark решает проблему на уровне архитектуры: промежуточные данные хранятся в памяти Executor-ов, а не записываются на диск после каждого шага.

Скорость разных типов хранения:

Тип Пропускная способность Latency
DDR4 RAM ~50 GB/s ~100 нс
NVMe SSD ~3–7 GB/s ~100 мкс
SATA SSD ~0.5 GB/s ~200 мкс
HDD (HDFS) ~0.1–0.2 GB/s ~10 мс

RAM быстрее HDD в 100–500 раз по latency и в 25–50 раз по пропускной способности. Именно это скрывается за формулой «10–100× быстрее».

Для итеративного алгоритма разница наглядна:

Если RAM не хватает - Spark сбрасывает (spill) часть данных на локальный диск Executor-а (не на HDFS), но это исключение. При достаточном объёме памяти итеративный алгоритм читает источник один раз, а дальнейшие итерации работают полностью в RAM.

Когда разница меньше

Важная оговорка: «10–100×» - это итеративные сценарии и сложные многошаговые запросы.

Сценарий Практический выигрыш
Итеративный ML (100 итераций) 50–100×
Сложный ETL: JOIN × 3, 10+ трансформаций 5–20×
Одиночный SELECT с GROUP BY 2–5×
Одиночное полное сканирование большого файла ~1–2×

Если задача - один проход по данным с записью на выход, Spark быстрее за счёт лучшего параллелизма и оптимизатора Catalyst, но не за счёт in-memory магии.


DAG: произвольный граф вместо Map → Reduce

MapReduce ограничен ровно двумя фазами. Сложный запрос требует нескольких последовательных Job-ов с HDFS между ними.

Spark строит DAG (Directed Acyclic Graph) - произвольный граф зависимостей произвольной глубины. Весь SQL-запрос с тремя JOIN-ами и несколькими агрегациями превращается в один DAG, который оптимизируется как единое целое.

Из урока 2 мы знаем: DAG разбивается на Stage-ы по границам shuffle. Внутри каждого Stage все операции выполняются в один проход по данным без промежуточных записей. На локальный диск пишутся только файлы shuffle - то, что нужно для передачи данных между Stage-ами. HDFS при этом не трогается.


Lazy Evaluation: откладываем всё до последнего

Это ключевая концепция урока. Spark не выполняет трансформации сразу. Когда вы пишете df.filter(...) или df.groupBy(...), Spark только записывает эту операцию в план - ничего не происходит.

Выполнение начинается только при вызове Action - операции, требующей реального результата.

# Ни одна из этих строк ничего не вычисляет.
# Spark читает схему Parquet-файла, но не данные.
df = spark.read.parquet("s3a://data/events/")

# Каждая строка ниже добавляет узел в DAG - не более.
filtered = df.filter(df.event_type == "purchase")
enriched = filtered.withColumn(
    "amount_usd", col("amount") / col("exchange_rate")
)
aggregated = enriched.groupBy("user_id") \
    .agg(sum("amount_usd").alias("total"))

# Вот здесь Spark начинает работать:
# читает Parquet, выполняет filter → withColumn → groupBy → agg → запись
aggregated.write.mode("overwrite").parquet("s3a://result/user_totals/")

Почему это важно? До Action у Spark есть весь план. Оптимизатор Catalyst видит все операции сразу и может перестроить их порядок, устранить лишние вычисления и передать условия фильтрации прямо в Parquet Reader - до того, как первый Task попадёт на Executor.

Если бы Spark выполнял каждую трансформацию немедленно, эти оптимизации были бы невозможны: к моменту write() про filter() уже было бы «забыто».

Трансформации и Actions - полный список

Тип Что делает Примеры
Transformation Добавляет операцию в DAG, возвращает новый DataFrame filter(), select(), groupBy(), join(), withColumn(), union(), repartition(), distinct()
Action Запускает выполнение DAG, возвращает результат count(), collect(), show(), first(), take(N), write.*(), foreach(), toPandas()

Практическое следствие: не бойтесь длинных цепочек трансформаций. Стоимость - не в их количестве, а в объёме данных и наличии shuffle. Две последовательных filter() - не два прохода, а один после оптимизации.


Catalyst оптимизирует ваш план до выполнения

Когда вы вызываете Action, Catalyst проходит несколько стадий оптимизации (детально разбираем в модуле об оптимизации). Два результата, которые ощущаются в производительности прямо сейчас:

Predicate Pushdown

Если в вашем коде есть filter(), Catalyst передаёт условие фильтрации вниз к хранилищу - Parquet Reader читает только строки, удовлетворяющие условию. Не поднимает весь файл в память, чтобы потом отфильтровать.

Parquet хранит статистику (min/max значения) для каждой группы строк (row group). Если min_date > '2024-01' - весь row group пропускается без чтения. Это не ошибка оптимизации - это работа Catalyst + формат Parquet вместе.

Column Pruning

SELECT user_id, amount из Parquet-файла с 80 колонками читает физически только 2 колонки. Catalyst передаёт список нужных колонок в Parquet Reader - остальные 78 не читаются вообще. В Avro (строковый формат) так не работает: придётся читать всю строку целиком.

Как проверить: explain()

explain() показывает физический план до выполнения. Вызывайте перед запуском на production:

df = spark.read.parquet("s3a://data/events/")

result = df.filter(col("event_type") == "purchase") \
           .select("user_id", "amount")

result.explain(mode="formatted")

Вывод (упрощённый):

== Physical Plan ==
*(1) Project [user_id#10L, amount#11]
+- *(1) Filter (isnotnull(event_type#12) AND (event_type#12 = purchase))
   +- FileScan parquet [user_id#10L, amount#11, event_type#12]
        Batched: true,
        DataFilters: [isnotnull(event_type#12), (event_type#12 = purchase)],
        Format: Parquet,
        PushedFilters: [IsNotNull(event_type), EqualTo(event_type,purchase)],
        ReadSchema: struct<user_id:bigint,amount:double,event_type:string>

Что читать:

  • PushedFilters - условия переданы в Parquet Reader. Если пусто - pushdown не сработал (проверьте тип колонки и формат).
  • ReadSchema - только 3 колонки из исходных 80. Column pruning работает.
  • FileScan - чтение с диска / S3, всегда в самом низу плана.

Для сложных запросов с join-ами используйте:

result.explain(mode="extended")  # показывает все стадии: Analyzed → Optimized → Physical

Кэширование: когда данные нужны повторно

Lazy Evaluation не всегда плюс. Если один и тот же DataFrame используется в нескольких Action-ах, Spark пересчитывает его с нуля каждый раз: перечитывает источник, повторяет все трансформации.

# ❌ Плохо: df_cleaned пересчитывается дважды, S3 читается дважды
df_cleaned = raw.filter(col("status") == "active") \
                .dropDuplicates(["user_id"])

row_count = df_cleaned.count()          # Action 1: читаем S3, считаем
df_cleaned.write.parquet("s3a://out/")  # Action 2: читаем S3 снова

Решение - persist():

# ✅ Хорошо: читаем источник один раз, результат остаётся в RAM
from pyspark import StorageLevel

df_cleaned = raw.filter(col("status") == "active") \
                .dropDuplicates(["user_id"]) \
                .persist(StorageLevel.MEMORY_AND_DISK)

row_count = df_cleaned.count()          # Action 1: читаем S3, результат → RAM
df_cleaned.write.parquet("s3a://out/")  # Action 2: читаем из RAM

df_cleaned.unpersist()                  # Явно освобождаем память

persist() срабатывает при первом Action - это тоже Lazy Evaluation. Второй Action уже читает из памяти.

cache() - алиас для persist(StorageLevel.MEMORY_AND_DISK_DESER). Разница в StorageLevel:

Уровень Где Когда использовать
MEMORY_AND_DISK RAM, переполнение → диск По умолчанию. Защита от OOM
MEMORY_ONLY Только RAM DataFrame небольшой, нужна максимальная скорость
DISK_ONLY Только диск RAM ограничена, но пересчёт дорогой (сложный join)

Когда кэшировать:

  • DataFrame используется в 2+ Action-ах
  • Вычисление дорогое: несколько join-ов, дедупликация, сложная агрегация
  • Итеративный алгоритм (MLlib, Graph-алгоритмы)

Когда не кэшировать:

  • DataFrame читается один раз - только лишний overhead на сериализацию
  • Данные не помещаются в RAM → постоянный spill, выигрыш нулевой
  • Простой проход: filter → select → write без повторного использования

Всегда вызывайте unpersist() явно: Spark не знает, когда вы закончили использовать кэш. Накопленные кэши вытесняют рабочую память и замедляют shuffle.


Антипаттерн: collect() убивает Driver

collect() переносит все данные с Executor-ов в память Driver-процесса. Driver - один JVM-процесс с фиксированным spark.driver.memory (по умолчанию 1–2 GB).

# ❌ Катастрофа: 10 млн строк × 10 колонок → Driver OOM → приложение падает
all_rows = df.collect()

# ✅ Вместо этого - пишите результат в хранилище
df.write.mode("overwrite").parquet("s3a://result/")

# ✅ Для отладки - берите ограниченную выборку
sample = df.limit(100).collect()   # только 100 строк безопасно

# ✅ Для агрегатных метрик - collect() безопасен на маленьком результате
stats = df.agg(
    count("*").alias("total"),
    sum("amount").alias("sum_amount"),
    avg("amount").alias("avg_amount")
).collect()
# Результат - список из одной Row с тремя числами
print(stats[0]["total"], stats[0]["sum_amount"])

Правило: collect() допустим только если вы точно знаете, что результат мал - после limit(), на агрегированном результате с малым числом ключей, или на groupBy() с заведомо ограниченным числом групп.

При падении по OOM в Spark UI увидите: ExecutorLostFailure или java.lang.OutOfMemoryError: Java heap space на Driver-узле.


Практика: наблюдаем Lazy Evaluation в Spark UI

Запустите код и откройте Spark UI (http://localhost:4040):

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, sum as _sum, count

spark = SparkSession.builder \
    .appName("LazyEvalDemo") \
    .master("local[*]") \
    .getOrCreate()

# Датасет в памяти для локального теста
data = [(i, i % 10, float(i * 1.5)) for i in range(1_000_000)]
df = spark.createDataFrame(data, ["id", "category", "value"])

# ── Трансформации: в Spark UI нет ни одного Job ──
step1 = df.filter(col("value") > 100.0)
step2 = step1.withColumn("value_doubled", col("value") * 2)
step3 = step2.filter(col("category") < 5)
step4 = step3.groupBy("category").agg(
    count("*").alias("cnt"),
    _sum("value").alias("total")
)

# Проверяем план ДО выполнения
step4.explain(mode="formatted")

# ── Action: появится один Job в Spark UI ──
result = step4.collect()
print(f"Результат: {len(result)} групп")

# В Spark UI → Jobs: один Job, два Stage
# Stage 1: filter + withColumn + filter (всё narrow, в один проход)
# Stage 2: groupBy + agg (wide - нужен shuffle)

Ключевое наблюдение: четыре отдельных шага (filter, withColumn, filter, groupBy) Spark выполнил в два Stage - потому что первые три narrow-трансформации склеились в один проход. Это Lazy Evaluation в действии: Spark знал план заранее и объединил операции.


Итог: Spark vs MapReduce

Аспект MapReduce Spark
Промежуточные данные HDFS (диск) RAM + spill на локальный диск
Граф зависимостей Map → Reduce (2 фазы) DAG произвольной глубины
Сложный запрос (JOIN × 3) 3 Job-а, 6 записей HDFS 1 DAG, 1–2 Stage
Итеративный ML (100 итер.) 200 операций HDFS 1 чтение + 99 итераций в RAM
Оптимизация Нет Catalyst: pushdown, pruning, AQE
Python API Hadoop Streaming (неудобно) PySpark - полноценный API

MapReduce сегодня практически не используется в новых проектах. Если в компании есть Hadoop-кластер, Spark запускается поверх него через YARN - аппаратная инфраструктура остаётся, вычислительный движок меняется.

Чек-лист

  • [ ] Трансформации не выполняются - каждая добавляет узел в DAG
  • [ ] Action запускает весь DAG целиком
  • [ ] explain() перед запуском на production - убедись, что видишь PushedFilters
  • [ ] persist() только если DataFrame используется в 2+ Action-ах
  • [ ] unpersist() явно после использования кэша
  • [ ] collect() только на заведомо маленьком результате (после limit() или на агрегате)
  • [ ] Для записи результатов - write.parquet(), не collect() + python-файл

В следующем уроке разберём RDD - базовую абстракцию Spark, из которой вырос DataFrame API, и поймём, почему в 2024 году она почти не используется напрямую, но знать её нужно.